3 Commits
Author SHA1 Message Date
sneak 695bf8665e Each target's row shows its result as soon as its check ends (closes #91)
check / check (push) Successful in 2m22s
tick drew no row until the round's slowest check ended, up to 24
seconds at a 30-second interval since checks time out at 80% of it.
Each check now pushes its sample and redraws its row as it ends. When
the last check ends, every row is redrawn, so none still reads "paused"
after a pause and resume, and sorting, the summary, the health box and
offline detection run once. A check that ends while paused or after its
round is given up draws nothing, and the first round is still discarded
as a whole. The row is looked up when the check ends, as a pin click can
re-sort the rows mid-round.

Unit tests run tick on the mocked clock against a stand-in page.

Model: opus-5-5
2026-10-03 23:32:05 +00:00
clawbot 39ee6ca839 Entrypoint acts as root on nothing outside /data (closes #80)
check / check (push) Successful in 1m58s
`bin/entrypoint.sh` now runs `netwatch-server prepare-data-dir`, which
refuses a `DATA_DIR` that is not `/data` or a path below it written in
full, then creates `DATA_DIR`, gives `/data` and everything in it to
`netwatch`, and sets mode 750 on `/data` and `DATA_DIR`. Every step goes
through a Go `os.Root` opened on `/data`, and the modes are set on the
opened directories rather than by name, so neither a symbolic link
already there nor one a host process swaps in while the container
starts can make root create or change anything outside `/data`. The
README says which `DATA_DIR` values are accepted.

Model: opus-5-5
2026-10-03 18:24:38 +02:00
clawbot e4df415676 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
2026-10-03 17:51:08 +02:00
12 changed files with 1139 additions and 132 deletions
+5 -3
View File
@@ -223,12 +223,14 @@ 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 - `REPORTS_PER_MINUTE`, default `60`: reports each client address may send a
minute minute
- `DATA_DIR_MAX_BYTES`, default `1073741824` (1 GiB): the most room the - `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 - `CORS_ALLOWED_ORIGINS`, default empty: other origins whose pages may call
the API the API
- `DEBUG`, default `false`: debug logging - `DEBUG`, default `false`: debug logging
- `DATA_DIR`, default `/data/reports`: leave unset; reports kept outside - `DATA_DIR`, default `/data/reports`: the directory the reports are kept
`/data` do not survive a redeploy in: `/data` or a path below it, with no `.` or `..` part and no extra `/`.
The container also stops if the path goes through a symbolic link that
leads out of `/data` or is written as a full path
- `TRUSTED_PROXIES`, default empty: set it to the address the reverse proxy - `TRUSTED_PROXIES`, default empty: set it to the address the reverse proxy
in front of the container connects from, as an IP address or CIDR; several in front of the container connects from, as an IP address or CIDR; several
are separated by commas. nginx takes the client address from are separated by commas. nginx takes the client address from
+19
View File
@@ -30,6 +30,25 @@ latest run passes.
round's last check ends, so no row reads "paused" after a pause and resume round's last check ends, so no row reads "paused" after a pause and resume
during the round. A check that ends after the user pauses or after its round during the round. A check that ends after the user pauses or after its round
is given up shows nothing, and the first round is still discarded as a whole is given up shows nothing, and the first round is still discarded as a whole
- 2026-10-03: root no longer acts outside `/data` when it prepares `DATA_DIR`
(issue #80): `bin/entrypoint.sh` runs `netwatch-server prepare-data-dir`,
which refuses a `DATA_DIR` that is not `/data` or a path below it written in
full, then creates `DATA_DIR`, gives `/data` and everything in it to
`netwatch` and sets the modes, all through a Go `os.Root` opened on `/data`.
That refuses any path leading out of `/data`, so neither a symbolic link
already there nor one a host process swaps in during the start can make root
create or change anything elsewhere, and `DATA_DIR=/etc` no longer gives
`/etc` to `netwatch`. The `README.md` section "Running under upaas" says which
values are accepted
- 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: `backend/script/lint` says what went wrong with its - 2026-10-03: `backend/script/lint` says what went wrong with its
`.golangci.yml` check (issue #34). On a hash mismatch it says to compare the `.golangci.yml` check (issue #34). On a hash mismatch it says to compare the
file with the org standard: if they differ, restore the org standard; if they file with the org standard: if they differ, restore the org standard; if they
+20 -13
View File
@@ -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 | | `BIND_ADDRESS` | empty | IP address to listen on; empty listens on every interface |
| `PORT` | `8080` | HTTP listen port | | `PORT` | `8080` | HTTP listen port |
| `DATA_DIR` | `./data/reports` | Directory for compressed reports | | `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 | | `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 | | `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) | | `REPORTS_PER_MINUTE` | `60` | Reports each client address may send a minute; see [Report limits](#report-limits) |
@@ -107,8 +107,9 @@ this server. The image's entrypoint, `bin/entrypoint.sh`, starts the server as
user `netwatch` (uid 1000) with `BIND_ADDRESS=127.0.0.1` and `PORT=8081`, so user `netwatch` (uid 1000) with `BIND_ADDRESS=127.0.0.1` and `PORT=8081`, so
only nginx reaches it, and with `TRUSTED_PROXIES=127.0.0.1/32`, so it takes the only nginx reaches it, and with `TRUSTED_PROXIES=127.0.0.1/32`, so it takes the
client address nginx passes on and no other. `DATA_DIR` is `/data/reports`, on client address nginx passes on and no other. `DATA_DIR` is `/data/reports`, on
the `/data` volume; the entrypoint creates it and gives it and `/data` to the `/data` volume; before starting the server, the entrypoint creates it and
`netwatch` before starting the server. nginx replaces the security headers gives it and `/data` to `netwatch` with `netwatch-server prepare-data-dir`,
which acts on nothing outside `/data`. nginx replaces the security headers
this server sets with those in the root `security-headers.conf`, so those are this server sets with those in the root `security-headers.conf`, so those are
what clients of the image see. what clients of the image see.
@@ -128,10 +129,11 @@ Reports are written as `reports-<timestamp>-<number>.jsonl.zst` files in
`DATA_DIR`. The timestamp is in UTC to the millisecond, so the names sort by `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 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 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 still get different names. A failed write uses up its number and leaves a gap in
the numbers if the file could not be created and otherwise a file under that the numbers: its file, if it was created, is removed. The file stays, counted
number that may be incomplete. Each file contains one JSON object per line, toward `DATA_DIR_MAX_BYTES` from the next start, only if removing it fails too.
compressed with zstd. Files are created with `O_EXCL` to prevent overwrites. Each file contains one JSON object per line, compressed with zstd. Files are
created with `O_EXCL` to prevent overwrites.
### Report limits ### Report limits
@@ -151,11 +153,17 @@ credentials, so it is bounded instead. Both refusals below answer with the same
`X-RateLimit-Reset` headers. `X-RateLimit-Reset` headers.
- **Size cap.** The report files in `DATA_DIR` may total at most - **Size cap.** The report files in `DATA_DIR` may total at most
`DATA_DIR_MAX_BYTES`, counting the files already there at start. Reports `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 waiting in memory count at their uncompressed size until they are written;
a report that would take the total past the cap is refused with 507, and those lost to a failed write stop counting, and the part of its file written
nothing of it is stored. Deleting report files frees room only at the next is removed. When a report would take the total past the cap, the oldest report
start, when the files are counted again. The default of 1 GiB is small enough files are deleted to make room, and each deletion is logged with the file's
for any host; set it to the space you can give `DATA_DIR`. 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 ### CORS
@@ -173,7 +181,6 @@ starting, with an error naming `CORS_ALLOWED_ORIGINS`.
- Add integration test that POSTs a report and verifies the compressed output - Add integration test that POSTs a report and verifies the compressed output
- Add report decompression/query endpoint - Add report decompression/query endpoint
- Add metrics (Prometheus) for buffer size, flush count, report count - Add metrics (Prometheus) for buffer size, flush count, report count
- Add retention policy to prune old report files
## License ## License
+21
View File
@@ -4,6 +4,7 @@ package main
import ( import (
"fmt" "fmt"
"os" "os"
"os/user"
"sneak.berlin/go/netwatch/internal/config" "sneak.berlin/go/netwatch/internal/config"
"sneak.berlin/go/netwatch/internal/globals" "sneak.berlin/go/netwatch/internal/globals"
@@ -37,6 +38,26 @@ func main() {
return return
} }
// "netwatch-server prepare-data-dir DATA_DIR" gets DATA_DIR ready
// for the netwatch user, or exits 1 with the error; see
// reportbuf.PrepareDataDir. bin/entrypoint.sh runs it as root
// before it starts this server as that user.
if len(os.Args) == 3 && os.Args[1] == "prepare-data-dir" {
netwatch, err := user.Lookup("netwatch")
if err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
err = reportbuf.PrepareDataDir("/data", os.Args[2], netwatch)
if err != nil {
fmt.Fprintf(os.Stderr, "DATA_DIR '%s': %v\n", os.Args[2], err)
os.Exit(1)
}
return
}
globals.Appname = Appname globals.Appname = Appname
globals.Version = Version globals.Version = Version
+4 -3
View File
@@ -86,11 +86,12 @@ func (s *Handlers) decodeErrorStatus(err error) int {
} }
// appendErrorStatus logs a failure to store a report and returns // appendErrorStatus logs a failure to store a report and returns
// the status to send: 507 when the report files are at their size // the status to send: 507 when the reports waiting to be written fill
// cap, otherwise 500. // the size cap, otherwise 500.
func (s *Handlers) appendErrorStatus(err error) int { func (s *Handlers) appendErrorStatus(err error) int {
if errors.Is(err, reportbuf.ErrFull) { 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 return http.StatusInsufficientStorage
} }
+3 -3
View File
@@ -125,9 +125,9 @@ func TestHandleReportStorageFailureIsNon2xx(t *testing.T) {
} }
} }
// TestHandleReportFullIs507 checks the answer when the report files // TestHandleReportFullIs507 checks the answer when the reports waiting
// are at their size cap: 507 and the usual error body, which tells // to be written fill the size cap: 507 and the usual error body, which
// the client nothing more. // tells the client nothing more.
func TestHandleReportFullIs507(t *testing.T) { func TestHandleReportFullIs507(t *testing.T) {
t.Parallel() t.Parallel()
+150
View File
@@ -0,0 +1,150 @@
package reportbuf
import (
"errors"
"io/fs"
"os"
"os/user"
"path/filepath"
"strconv"
"syscall"
)
// ErrDataDirOutsideVolume is returned by PrepareDataDir for a DATA_DIR
// that is not the volume or a path below it, written in full.
var ErrDataDirOutsideVolume = errors.New(
"must be /data or a path below it, with no '.', '..' or extra '/'")
// PrepareDataDir gets dir, the server's DATA_DIR, ready for owner, the
// user the server runs as, so that a host directory mounted at volume,
// /data in the image, needs no preparing: it creates dir, gives volume
// and everything in it to owner, and gives volume and dir the mode the
// server gives a directory it creates. dir must be volume or a path
// below it, with no '.', '..', empty part or '/' at the end.
//
// bin/entrypoint.sh runs this as root, which would follow a symbolic
// link anywhere, so every step goes through an os.Root opened on
// volume: it follows a link only when it is written as a relative
// path that stays inside volume, and refuses any other. A process on
// the host can swap a link onto a path in volume at any moment while
// this runs. Even then, the os.Root checks each link as it reaches
// it. MkdirAll creates each directory inside a parent it already has
// open, never following a link at the name it creates, and follows a
// link on the path only as the os.Root allows, so a relative link
// inside volume can lead it to create directories elsewhere inside
// volume. Lchown never changes what a link points to, and the modes
// are set on directories already opened (see chmodDir), so the most
// that process can do is make a step fail or wait, or act on
// something else inside volume.
func PrepareDataDir(volume, dir string, owner *user.User) error {
// rel is dir as a path from volume; IsLocal is false for one that
// leads out of it.
rel, err := filepath.Rel(volume, dir)
if err != nil || dir != filepath.Clean(dir) || !filepath.IsLocal(rel) {
return ErrDataDirOutsideVolume
}
uid, err := strconv.Atoi(owner.Uid)
if err != nil {
return err
}
gid, err := strconv.Atoi(owner.Gid)
if err != nil {
return err
}
root, err := os.OpenRoot(volume)
if err != nil {
return err
}
defer func() { _ = root.Close() }()
err = root.MkdirAll(rel, dirPerms)
if err != nil {
return err
}
err = lchownAll(root, ".", uid, gid)
if err != nil {
return err
}
err = chmodDir(root, ".")
if err != nil {
return err
}
return chmodDir(root, rel)
}
// lchownAll gives name, a directory inside root, and everything in it
// to uid and gid. It reads each directory opened through root, not
// through root.FS(), which refuses a name that is not valid UTF-8, and
// calls Lchown on every entry, which gives a symbolic link itself to
// them, not what it points to. It goes into an entry only when the
// read found a directory there, so it follows no link it finds; one
// swapped in for that directory afterwards is followed only as the
// os.Root allows.
func lchownAll(root *os.Root, name string, uid, gid int) error {
err := root.Lchown(name, uid, gid)
if err != nil {
return err
}
dir, err := root.Open(name)
if err != nil {
return err
}
entries, err := dir.ReadDir(-1)
_ = dir.Close()
if err != nil {
return err
}
for _, entry := range entries {
entryName := filepath.Join(name, entry.Name())
if entry.IsDir() {
err = lchownAll(root, entryName, uid, gid)
} else {
err = root.Lchown(entryName, uid, gid)
}
if err != nil {
return err
}
}
return nil
}
// chmodDir gives name, a directory inside root, the mode the server
// gives a directory it creates. Root.Chmod would not hold: on Linux it
// checks that name is not a symbolic link, then sets the mode by name,
// following a link swapped in between. So chmodDir opens name through
// root and sets the mode on the open directory. It refuses anything
// but a directory: a directory has no second name (hard link), so the
// one opened is inside root, where any other file could be a hard link
// to one outside.
func chmodDir(root *os.Root, name string) error {
dir, err := root.Open(name)
if err != nil {
return err
}
defer func() { _ = dir.Close() }()
info, err := dir.Stat()
if err != nil {
return err
}
if !info.IsDir() {
return &fs.PathError{Op: "chmod", Path: name, Err: syscall.ENOTDIR}
}
return dir.Chmod(dirPerms)
}
+307
View File
@@ -0,0 +1,307 @@
package reportbuf_test
import (
"errors"
"io/fs"
"os"
"os/user"
"path/filepath"
"strconv"
"syscall"
"testing"
"sneak.berlin/go/netwatch/internal/reportbuf"
)
// reports is the last part of DATA_DIR in these tests, as in the
// image's /data/reports.
const reports = "reports"
// currentUser is the user the test runs as, the only owner a test not
// run as root can give files to.
func currentUser() *user.User {
return &user.User{
Uid: strconv.Itoa(os.Getuid()),
Gid: strconv.Itoa(os.Getgid()),
}
}
// tempDirMode700 is a new directory in a t.TempDir with mode 0700, so
// a test can tell that PrepareDataDir left its mode alone.
func tempDirMode700(t *testing.T) string {
t.Helper()
dir := filepath.Join(t.TempDir(), "d")
err := os.Mkdir(dir, 0o700)
if err != nil {
t.Fatal(err)
}
return dir
}
func requireMode(t *testing.T, path string, want fs.FileMode) {
t.Helper()
info, err := os.Stat(path)
if err != nil {
t.Fatal(err)
}
if info.Mode() != want {
t.Errorf("%s: mode %v, want %v", path, info.Mode(), want)
}
}
func requireMissing(t *testing.T, path string) {
t.Helper()
_, err := os.Lstat(path)
if !errors.Is(err, fs.ErrNotExist) {
t.Errorf("%s: created, or Lstat failed: %v", path, err)
}
}
func requireOwner(t *testing.T, path string, uid, gid uint32) {
t.Helper()
info, err := os.Lstat(path)
if err != nil {
t.Fatal(err)
}
stat, _ := info.Sys().(*syscall.Stat_t)
if stat.Uid != uid || stat.Gid != gid {
t.Errorf("%s: owner %d:%d, want %d:%d", path, stat.Uid, stat.Gid,
uid, gid)
}
}
func TestPrepareDataDirCreatesDataDir(t *testing.T) {
t.Parallel()
volume := t.TempDir()
dir := filepath.Join(volume, "a", reports)
err := reportbuf.PrepareDataDir(volume, dir, currentUser())
if err != nil {
t.Fatal(err)
}
requireMode(t, volume, fs.ModeDir|0o750)
requireMode(t, dir, fs.ModeDir|0o750)
}
// TestPrepareDataDirSetsModeOfExistingDataDir: a DATA_DIR already on
// the host with another mode gets the mode too, not only a new one.
func TestPrepareDataDirSetsModeOfExistingDataDir(t *testing.T) {
t.Parallel()
volume := t.TempDir()
dir := filepath.Join(volume, reports)
err := os.Mkdir(dir, 0o700)
if err != nil {
t.Fatal(err)
}
err = reportbuf.PrepareDataDir(volume, dir, currentUser())
if err != nil {
t.Fatal(err)
}
requireMode(t, dir, fs.ModeDir|0o750)
}
func TestPrepareDataDirTakesTheVolumeItself(t *testing.T) {
t.Parallel()
volume := t.TempDir()
err := reportbuf.PrepareDataDir(volume, volume, currentUser())
if err != nil {
t.Fatal(err)
}
requireMode(t, volume, fs.ModeDir|0o750)
}
// TestPrepareDataDirRefusesDataDirOutsideVolume covers a DATA_DIR that
// is relative, outside the volume, or not written in full.
func TestPrepareDataDirRefusesDataDirOutsideVolume(t *testing.T) {
t.Parallel()
volume := tempDirMode700(t)
for _, dir := range []string{
reports, volume + "/../new", volume + "/", volume + "//" + reports,
volume + "/./" + reports, volume + "/" + reports + "/..", volume + "x",
"/etc",
} {
err := reportbuf.PrepareDataDir(volume, dir, currentUser())
if !errors.Is(err, reportbuf.ErrDataDirOutsideVolume) {
t.Errorf("%q: error = %v, want ErrDataDirOutsideVolume", dir, err)
}
}
requireMissing(t, filepath.Join(filepath.Dir(volume), "new"))
requireMode(t, volume, fs.ModeDir|0o700)
}
// TestPrepareDataDirRefusesLinkOutOfVolume puts a symbolic link to a
// directory outside the volume on the path to DATA_DIR, written as a
// full path and as one that climbs out with '..', and as DATA_DIR
// itself, where the mode would be set through it.
func TestPrepareDataDirRefusesLinkOutOfVolume(t *testing.T) {
t.Parallel()
for _, tc := range []struct {
name string
climbsOut bool
link, dir string
}{
{"full path", false, "x", "x/reports"},
{"climbs out", true, "x", "x/reports"},
{"DATA_DIR itself", false, reports, reports},
} {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
outside := tempDirMode700(t)
volume := t.TempDir()
climbOut, err := filepath.Rel(volume, outside)
if err != nil {
t.Fatal(err)
}
target := outside
if tc.climbsOut {
target = climbOut
}
err = os.Symlink(target, filepath.Join(volume, tc.link))
if err != nil {
t.Fatal(err)
}
err = reportbuf.PrepareDataDir(volume,
filepath.Join(volume, tc.dir), currentUser())
if err == nil {
t.Error("no error")
}
requireMissing(t, filepath.Join(outside, reports))
requireMode(t, outside, fs.ModeDir|0o700)
})
}
}
// TestPrepareDataDirRefusesDanglingLink: DATA_DIR is a symbolic link
// to a name in the volume that does not exist, which is not created.
func TestPrepareDataDirRefusesDanglingLink(t *testing.T) {
t.Parallel()
volume := t.TempDir()
err := os.Symlink("missing", filepath.Join(volume, reports))
if err != nil {
t.Fatal(err)
}
err = reportbuf.PrepareDataDir(volume, filepath.Join(volume, reports),
currentUser())
if err == nil {
t.Error("no error")
}
requireMissing(t, filepath.Join(volume, "missing"))
}
// TestPrepareDataDirTakesDirectoryNamedInLatin1: a host directory can
// hold names that are not valid UTF-8, here "café" written in Latin-1.
// A test not run as root can only check that PrepareDataDir goes into
// such a directory and gives it, and what it holds, to the current
// user.
func TestPrepareDataDirTakesDirectoryNamedInLatin1(t *testing.T) {
t.Parallel()
volume := t.TempDir()
latin1 := filepath.Join(volume, "caf\xe9")
err := os.Mkdir(latin1, 0o700)
if err != nil {
t.Fatal(err)
}
err = os.WriteFile(filepath.Join(latin1, "f"), nil, 0o600)
if err != nil {
t.Fatal(err)
}
owner := currentUser()
err = reportbuf.PrepareDataDir(volume, filepath.Join(volume, reports),
owner)
if err != nil {
t.Fatal(err)
}
uid, _ := strconv.ParseUint(owner.Uid, 10, 32)
gid, _ := strconv.ParseUint(owner.Gid, 10, 32)
requireOwner(t, latin1, uint32(uid), uint32(gid))
requireOwner(t, filepath.Join(latin1, "f"), uint32(uid), uint32(gid))
}
// TestPrepareDataDirGivesVolumeToOwner gives everything in the volume
// to a uid and gid that own nothing, which only root can do. A
// symbolic link in the volume to a directory outside it is given to
// them itself; what it points to is left as it was.
func TestPrepareDataDirGivesVolumeToOwner(t *testing.T) {
t.Parallel()
if os.Geteuid() != 0 {
t.Skip("only root can give files to another uid")
}
outside := t.TempDir()
volume := t.TempDir()
old := filepath.Join(volume, "old")
err := os.WriteFile(filepath.Join(outside, "f"), nil, 0o600)
if err != nil {
t.Fatal(err)
}
err = os.Mkdir(old, 0o700)
if err != nil {
t.Fatal(err)
}
err = os.WriteFile(filepath.Join(old, "f"), nil, 0o600)
if err != nil {
t.Fatal(err)
}
err = os.Symlink(outside, filepath.Join(old, "link"))
if err != nil {
t.Fatal(err)
}
err = reportbuf.PrepareDataDir(volume, filepath.Join(volume, reports),
&user.User{Uid: "4242", Gid: "4343"})
if err != nil {
t.Fatal(err)
}
for _, path := range []string{
volume, filepath.Join(volume, reports), old,
filepath.Join(old, "f"), filepath.Join(old, "link"),
} {
requireOwner(t, path, 4242, 4343)
}
requireOwner(t, outside, 0, 0)
requireOwner(t, filepath.Join(outside, "f"), 0, 0)
}
+11 -1
View File
@@ -1,6 +1,9 @@
package reportbuf package reportbuf
import "time" import (
"os"
"time"
)
// FlushSizeThreshold exposes the buffer size at which Append starts // FlushSizeThreshold exposes the buffer size at which Append starts
// writing a report file to the external tests. // writing a report file to the external tests.
@@ -17,3 +20,10 @@ func (b *Buffer) Flush() error {
func (b *Buffer) StopClock(at time.Time) { func (b *Buffer) StopClock(at time.Time) {
b.now = func() time.Time { return at } 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 // 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 package reportbuf
import ( import (
@@ -12,6 +13,7 @@ import (
"log/slog" "log/slog"
"os" "os"
"path/filepath" "path/filepath"
"slices"
"strings" "strings"
"sync" "sync"
"sync/atomic" "sync/atomic"
@@ -37,9 +39,10 @@ const (
fileSuffix = ".jsonl.zst" fileSuffix = ".jsonl.zst"
) )
// ErrFull is returned by Append when storing the report would // ErrFull is returned by Append when the reports waiting to be
// take the report files past the configured maximum size. // written leave no room for the report under the configured maximum
var ErrFull = errors.New("report files at their size cap") // 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. // Params defines the dependencies for Buffer.
type Params struct { type Params struct {
@@ -49,15 +52,34 @@ type Params struct {
Logger *logger.Logger 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 // Buffer accumulates JSON lines in memory and flushes them
// to zstd-compressed files on disk. // to zstd-compressed files on disk.
type Buffer struct { type Buffer struct {
buf bytes.Buffer buf bytes.Buffer
dataDir string dataDir string
done chan struct{} done chan struct{}
log *slog.Logger // fileCreated is called with each report file once it is
maxBytes int64 // created, before anything is written to it: it does nothing,
mu sync.Mutex // 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 // now is the clock report files are named by: time.Now, except
// in tests that need two flushes to share a timestamp. // in tests that need two flushes to share a timestamp.
now func() time.Time now func() time.Time
@@ -66,8 +88,8 @@ type Buffer struct {
seq atomic.Uint64 seq atomic.Uint64
stopOnce sync.Once stopOnce sync.Once
// usedBytes is what Append checks against maxBytes: the size // usedBytes is what Append checks against maxBytes: the size
// of the report files in dataDir, plus the reports not yet // of the report files in dataDir, plus the reports waiting to
// written to one at their uncompressed size. // be written to one at their uncompressed size.
usedBytes int64 usedBytes int64
} }
@@ -83,11 +105,12 @@ func New(
} }
b := &Buffer{ b := &Buffer{
dataDir: dir, dataDir: dir,
done: make(chan struct{}), done: make(chan struct{}),
log: params.Logger.Get(), fileCreated: func(*os.File) {},
maxBytes: params.Config.DataDirMaxBytes, log: params.Logger.Get(),
now: time.Now, maxBytes: params.Config.DataDirMaxBytes,
now: time.Now,
} }
lc.Append(fx.Hook{ lc.Append(fx.Hook{
@@ -97,12 +120,28 @@ func New(
return fmt.Errorf("create data dir: %w", err) return fmt.Errorf("create data dir: %w", err)
} }
// Report files left by earlier runs count too. // Report files left by earlier runs count too, and are
b.usedBytes, err = reportFilesSize(b.dataDir) // the first to be deleted to make room.
files, err := reportFiles(b.dataDir)
if err != nil { if err != nil {
return err 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() go b.flushLoop()
return nil return nil
@@ -128,9 +167,10 @@ func New(
} }
// Append marshals v as a single JSON line and appends it to // Append marshals v as a single JSON line and appends it to
// the buffer. It stores nothing and returns ErrFull if the line // the buffer. If the line would take usedBytes past maxBytes, the
// would take usedBytes past maxBytes. If the buffer reaches the // oldest report files are deleted to make room; it stores nothing
// size threshold, it is drained and written to disk // and returns ErrFull if that cannot make room. If the buffer
// reaches the size threshold, it is drained and written to disk
// asynchronously. // asynchronously.
func (b *Buffer) Append(v any) error { func (b *Buffer) Append(v any) error {
line, err := json.Marshal(v) line, err := json.Marshal(v)
@@ -142,6 +182,8 @@ func (b *Buffer) Append(v any) error {
b.mu.Lock() b.mu.Lock()
b.deleteOldestFiles(lineBytes)
if b.usedBytes+lineBytes > b.maxBytes { if b.usedBytes+lineBytes > b.maxBytes {
b.mu.Unlock() b.mu.Unlock()
@@ -171,6 +213,37 @@ func (b *Buffer) Append(v any) error {
return nil 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 // flushLoop runs a ticker that periodically flushes buffered
// data to disk until the done channel is closed. // data to disk until the done channel is closed.
func (b *Buffer) flushLoop() { func (b *Buffer) flushLoop() {
@@ -217,15 +290,48 @@ func (b *Buffer) drainBuf() []byte {
return data return data
} }
// writeFile creates a timestamped zstd-compressed JSONL file // writeFile writes data, reports drained from the buffer, to a new
// in the data directory. // timestamped zstd-compressed JSONL file in the data directory.
func (b *Buffer) writeFile(data []byte) error { func (b *Buffer) writeFile(data []byte) error {
// The timestamp comes first, so the names sort by time; the number // The timestamp comes first, so the names sort by time; the number
// after it tells apart files named in the same millisecond. // after it tells apart files named in the same millisecond.
ts := b.now().UTC().Format("2006-01-02T15-04-05.000Z") 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) 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 // path is built from the operator-supplied dataDir plus a
// generated timestamp and number, so it carries no external input. // generated timestamp and number, so it carries no external input.
f, err := os.OpenFile( //nolint:gosec // see comment above f, err := os.OpenFile( //nolint:gosec // see comment above
@@ -234,61 +340,65 @@ func (b *Buffer) writeFile(data []byte) error {
filePerms, filePerms,
) )
if err != nil { 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 b.fileCreated(f)
// path closes it explicitly to check the error; closing it
// a second time here is harmless.
defer func() { _ = f.Close() }()
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) enc, err := zstd.NewWriter(f)
if err != nil { if err != nil {
return fmt.Errorf("create zstd encoder: %w", err) return 0, fmt.Errorf("create zstd encoder: %w", err)
} }
_, err = enc.Write(data) _, err = enc.Write(data)
if err != nil { if err != nil {
_ = enc.Close() _ = enc.Close()
return fmt.Errorf("write compressed data: %w", err) return 0, fmt.Errorf("write compressed data: %w", err)
} }
err = enc.Close() err = enc.Close()
if err != nil { if err != nil {
return fmt.Errorf("close zstd encoder: %w", err) return 0, fmt.Errorf("close zstd encoder: %w", err)
} }
info, err := f.Stat() info, err := f.Stat()
if err != nil { if err != nil {
return fmt.Errorf("stat report file: %w", err) return 0, fmt.Errorf("stat report file: %w", err)
} }
err = f.Close() return info.Size(), nil
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
} }
// reportFilesSize returns the total size of the report files in // reportFiles returns the report files in dir, oldest first:
// dir. // os.ReadDir sorts them by name, and the names sort by time.
func reportFilesSize(dir string) (int64, error) { func reportFiles(dir string) ([]reportFile, error) {
entries, err := os.ReadDir(dir) entries, err := os.ReadDir(dir)
if err != nil { 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 { for _, entry := range entries {
name := entry.Name() name := entry.Name()
@@ -299,11 +409,11 @@ func reportFilesSize(dir string) (int64, error) {
info, err := entry.Info() info, err := entry.Info()
if err != nil { 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 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) { func TestAppendPastCapIsRefused(t *testing.T) {
report := map[string]string{"id": "cap"} report := map[string]string{"id": "cap"}
dir := t.TempDir()
earlier := reportFilePath(dir, 1)
t.Setenv("DATA_DIR", t.TempDir()) writeBytes(t, earlier, 1)
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(lineBytes(t, report)))
t.Setenv("DATA_DIR", dir)
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(1+lineBytes(t, report)))
buf := startBuffer(t) buf := startBuffer(t)
@@ -156,24 +164,38 @@ func TestAppendPastCapIsRefused(t *testing.T) {
if !errors.Is(err, reportbuf.ErrFull) { if !errors.Is(err, reportbuf.ErrFull) {
t.Fatalf("report past the cap: error = %v, want ErrFull", err) 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 // TestOldestReportFileDeletedFirst starts on a data directory holding
// directory holding a report file from an earlier run, and a file // report files from an earlier run, and a file that is not a report,
// that is not a report, which must not count. // which neither counts nor is ever deleted. Nothing is deleted while
func TestCapCountsReportFilesAlreadyInDataDir(t *testing.T) { // there is room; then only the oldest report file is.
const earlierBytes = 100 func TestOldestReportFileDeletedFirst(t *testing.T) {
const fileBytes = 100
report := map[string]string{"id": "cap"} report := map[string]string{"id": "oldest"}
dir := t.TempDir() 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"), writeBytes(t, oldest, fileBytes)
earlierBytes)
writeBytes(t, filepath.Join(dir, "notes.txt"), 10*earlierBytes) for _, path := range kept {
writeBytes(t, path, fileBytes)
}
t.Setenv("DATA_DIR", dir) t.Setenv("DATA_DIR", dir)
// Room for the three report files and one report.
t.Setenv("DATA_DIR_MAX_BYTES", t.Setenv("DATA_DIR_MAX_BYTES",
strconv.Itoa(earlierBytes+lineBytes(t, report))) strconv.Itoa(3*fileBytes+lineBytes(t, report)))
buf := startBuffer(t) buf := startBuffer(t)
@@ -182,9 +204,55 @@ func TestCapCountsReportFilesAlreadyInDataDir(t *testing.T) {
t.Fatalf("report that fills the cap exactly: %v", err) 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) err = buf.Append(report)
if !errors.Is(err, reportbuf.ErrFull) { if err != nil {
t.Fatalf("report past the cap: error = %v, want ErrFull", err) 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. // Repetitive, so its file is far smaller than its JSON.
report := map[string]string{"id": strings.Repeat("a", 1000)} report := map[string]string{"id": strings.Repeat("a", 1000)}
size := lineBytes(t, report) 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 // Room for the report twice over only if the first one counts
// at its file's size by the time the second arrives. // at its file's size by the time the second arrives.
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(2*size-1)) t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(2*size-1))
@@ -217,13 +286,22 @@ func TestWrittenReportsCountAtFileSize(t *testing.T) {
if err != nil { if err != nil {
t.Fatalf("second report, after the first was written: %v", err) 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 // TestWrittenFilesDeletedToMakeRoom writes one report file after
// under a small cap: each report must be taken while the files on disk // another under a small cap. Every report must be taken; files are
// leave room for it, and refused once they do not. // deleted only when the report would not fit beside them, and the
func TestWrittenReportsKeepCounting(t *testing.T) { // files kept leave room for it.
const maxBytes = 200 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"} report := map[string]string{"id": "written"}
size := int64(lineBytes(t, report)) size := int64(lineBytes(t, report))
@@ -234,23 +312,24 @@ func TestWrittenReportsKeepCounting(t *testing.T) {
buf := startBuffer(t) buf := startBuffer(t)
// Every file takes at least a byte, so they fill the cap within for range reports {
// maxBytes rounds. before := reportFilesBytes(t, dir)
for range maxBytes {
used := reportFilesBytes(t, dir)
err := buf.Append(report) err := buf.Append(report)
if used+size > maxBytes { if err != nil {
if !errors.Is(err, reportbuf.ErrFull) { t.Fatalf("with %d bytes of report files: %v", before, err)
t.Fatalf("with %d bytes of report files: error = %v, "+
"want ErrFull", used, err)
}
return
} }
if err != nil { after := reportFilesBytes(t, dir)
t.Fatalf("with %d bytes of report files: %v", used, err)
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() err = buf.Flush()
@@ -258,8 +337,217 @@ func TestWrittenReportsKeepCounting(t *testing.T) {
t.Fatalf("flush: %v", err) 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 // TestConcurrentAppendsStopAtCap appends from many goroutines at once
@@ -492,6 +780,14 @@ func readReportFiles(dir string) ([]string, error) {
return contents, nil 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) { func writeBytes(t *testing.T, path string, n int) {
t.Helper() 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 { func hasReportFile(t *testing.T, dir string) bool {
t.Helper() t.Helper()
@@ -524,3 +835,85 @@ func hasReportFile(t *testing.T, dir string) bool {
return false 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")
}
}
+6 -19
View File
@@ -63,26 +63,13 @@ done > /etc/nginx/trusted-proxies.conf
# netwatch-server keeps its report files in DATA_DIR, on the /data # netwatch-server keeps its report files in DATA_DIR, on the /data
# volume, which may be a host directory owned by root or by another # volume, which may be a host directory owned by root or by another
# uid. Both are given to the netwatch user here, with the mode the # uid. Here, as root, netwatch-server prepare-data-dir creates DATA_DIR
# server gives a directory it creates, so the host directory needs no # and gives /data and everything in it to the netwatch user, so the
# preparing. # host directory needs no preparing. It stops the start, naming
# # DATA_DIR, unless DATA_DIR is /data or a path below it, and it acts on
# chown and chmod, run as root, change whatever a symbolic link on the # nothing outside /data, whatever symbolic links it meets there.
# path points to, anywhere in the container, and the netwatch user can
# put one in /data. So the start stops unless readlink -f, which
# follows every link on a path, gives /data and DATA_DIR back as they
# are. It also writes a path in full, so a DATA_DIR with '.', '..' or
# an extra '/' in it is refused too.
export DATA_DIR="${DATA_DIR:-/data/reports}" export DATA_DIR="${DATA_DIR:-/data/reports}"
mkdir -p "$DATA_DIR" || exit 1 netwatch-server prepare-data-dir "$DATA_DIR" || exit 1
if [ "$(readlink -f /data)" != /data ] ||
[ "$(readlink -f "$DATA_DIR")" != "$DATA_DIR" ]; then
echo "entrypoint: DATA_DIR must be a full path with no '.', '..'," \
"extra '/' or symbolic link on it or on /data, not '$DATA_DIR'" >&2
exit 1
fi
chown -R netwatch:netwatch /data "$DATA_DIR" || exit 1
chmod 750 /data "$DATA_DIR" || exit 1
# A stop signal is only noted here; the loop below acts on it. # A stop signal is only noted here; the loop below acts on it.
stop_requested="" stop_requested=""