Compare commits
3
Commits
main
..
ee959fc33a
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ee959fc33a | ||
|
|
4ce0814b14 | ||
|
|
e6d6815ecb |
@@ -4,3 +4,6 @@ tmp
|
|||||||
.DS_Store
|
.DS_Store
|
||||||
*.log
|
*.log
|
||||||
.claude
|
.claude
|
||||||
|
|
||||||
|
# .git is sent so the build can stamp the version, without its config.
|
||||||
|
.git/config
|
||||||
|
|||||||
+17
-4
@@ -20,7 +20,7 @@ RUN make lint
|
|||||||
# golang:1.25-alpine (2026-02-27)
|
# golang:1.25-alpine (2026-02-27)
|
||||||
FROM golang:1.25-alpine@sha256:f6751d823c26342f9506c03797d2527668d095b0a15f1862cddb4d927a7a4ced AS builder
|
FROM golang:1.25-alpine@sha256:f6751d823c26342f9506c03797d2527668d095b0a15f1862cddb4d927a7a4ced AS builder
|
||||||
|
|
||||||
RUN apk add --no-cache make
|
RUN apk add --no-cache git make
|
||||||
|
|
||||||
WORKDIR /src
|
WORKDIR /src
|
||||||
|
|
||||||
@@ -37,11 +37,24 @@ RUN make test
|
|||||||
|
|
||||||
# make build is a shim around backend/script/build, the one definition
|
# make build is a shim around backend/script/build, the one definition
|
||||||
# of the build command:
|
# of the build command:
|
||||||
# CGO_ENABLED=0 go build -trimpath -ldflags "-s -w -X main.Version=... -X main.Buildarch=..."
|
# CGO_ENABLED=0 go build -trimpath -ldflags "-s -w -X main.Version=..."
|
||||||
# That script reads VERSION from the environment, so it is handed over
|
# That script reads VERSION from the environment, so it is handed over
|
||||||
# there rather than as a make variable.
|
# there rather than as a make variable.
|
||||||
ARG VERSION=dev
|
#
|
||||||
RUN VERSION="${VERSION}" make build
|
# The version is the VERSION build argument when one is given, otherwise
|
||||||
|
# `git describe --tags --always` of the repo's .git: the tag on a tagged
|
||||||
|
# commit, tag-N-gHASH on a commit after one, the short commit when no
|
||||||
|
# tag is reachable. A version that still comes out empty, dev or unknown
|
||||||
|
# fails the build. .git goes to /git, not /src/.git, where go build would
|
||||||
|
# find it and record VCS details of a work tree holding only backend/.
|
||||||
|
COPY .git /git
|
||||||
|
ARG VERSION
|
||||||
|
RUN version="${VERSION:-$(git --git-dir=/git describe --tags --always)}"; \
|
||||||
|
case "$version" in ""|dev|unknown) \
|
||||||
|
echo "version is '$version' although .git is present" >&2; \
|
||||||
|
exit 1 ;; \
|
||||||
|
esac; \
|
||||||
|
VERSION="$version" make build
|
||||||
|
|
||||||
# Frontend stage
|
# Frontend stage
|
||||||
# node:22-alpine as of 2026-02-22
|
# node:22-alpine as of 2026-02-22
|
||||||
|
|||||||
@@ -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
|
- `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
|
||||||
|
|||||||
@@ -23,6 +23,14 @@ latest run passes.
|
|||||||
|
|
||||||
# Completed Steps
|
# 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):
|
- 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
|
`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
|
and `/data` to the `netwatch` user with mode 750 before starting the backend
|
||||||
|
|||||||
+11
-8
@@ -32,7 +32,7 @@ pattern as the repo root: the targets in `backend/Makefile` are thin shims over
|
|||||||
`test`, `fmt` and `fmt-check`:
|
`test`, `fmt` and `fmt-check`:
|
||||||
|
|
||||||
- `script/build` — compile the static `netwatch-server` binary with its version
|
- `script/build` — compile the static `netwatch-server` binary with its version
|
||||||
and architecture stamped in. The version is `VERSION` from the environment;
|
stamped in. The version is `VERSION` from the environment;
|
||||||
when that is unset or empty, it falls back to `git describe` inside a git
|
when that is unset or empty, it falls back to `git describe` inside a git
|
||||||
checkout, then to `dev`
|
checkout, then to `dev`
|
||||||
- `script/test` — run the Go tests under a 30-second timeout
|
- `script/test` — run the Go tests under a 30-second timeout
|
||||||
@@ -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 |
|
| `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) |
|
||||||
@@ -148,11 +148,15 @@ 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
|
When a report would take the total past the cap, the oldest report files are
|
||||||
nothing of it is stored. Deleting report files frees room only at the next
|
deleted to make room, and each deletion is logged with the file's name and
|
||||||
start, when the files are counted again. The default of 1 GiB is small enough
|
size; a file still being written is never deleted. A report is refused with
|
||||||
for any host; set it to the space you can give `DATA_DIR`.
|
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
|
### 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 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
|
||||||
|
|
||||||
|
|||||||
@@ -19,9 +19,8 @@ import (
|
|||||||
|
|
||||||
//nolint:gochecknoglobals // set via ldflags at build time
|
//nolint:gochecknoglobals // set via ldflags at build time
|
||||||
var (
|
var (
|
||||||
Appname = "netwatch-server"
|
Appname = "netwatch-server"
|
||||||
Version string
|
Version string
|
||||||
Buildarch string
|
|
||||||
)
|
)
|
||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
@@ -40,7 +39,6 @@ func main() {
|
|||||||
|
|
||||||
globals.Appname = Appname
|
globals.Appname = Appname
|
||||||
globals.Version = Version
|
globals.Version = Version
|
||||||
globals.Buildarch = Buildarch
|
|
||||||
|
|
||||||
fx.New(
|
fx.New(
|
||||||
fx.Provide(
|
fx.Provide(
|
||||||
|
|||||||
@@ -10,22 +10,18 @@ var (
|
|||||||
Appname string
|
Appname string
|
||||||
// Version is the git version tag.
|
// Version is the git version tag.
|
||||||
Version string
|
Version string
|
||||||
// Buildarch is the build architecture.
|
|
||||||
Buildarch string
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// Globals holds build-time metadata for the application.
|
// Globals holds build-time metadata for the application.
|
||||||
type Globals struct {
|
type Globals struct {
|
||||||
Appname string
|
Appname string
|
||||||
Version string
|
Version string
|
||||||
Buildarch string
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// New creates a Globals instance from package-level variables.
|
// New creates a Globals instance from package-level variables.
|
||||||
func New(_ fx.Lifecycle) (*Globals, error) {
|
func New(_ fx.Lifecycle) (*Globals, error) {
|
||||||
return &Globals{
|
return &Globals{
|
||||||
Appname: Appname,
|
Appname: Appname,
|
||||||
Buildarch: Buildarch,
|
Version: Version,
|
||||||
Version: Version,
|
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ package logger
|
|||||||
import (
|
import (
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"os"
|
"os"
|
||||||
|
"runtime"
|
||||||
|
|
||||||
"sneak.berlin/go/netwatch/internal/globals"
|
"sneak.berlin/go/netwatch/internal/globals"
|
||||||
|
|
||||||
@@ -95,6 +96,6 @@ func (l *Logger) Identify() {
|
|||||||
l.log.Info("starting",
|
l.log.Info("starting",
|
||||||
"appname", l.params.Globals.Appname,
|
"appname", l.params.Globals.Appname,
|
||||||
"version", l.params.Globals.Version,
|
"version", l.params.Globals.Version,
|
||||||
"buildarch", l.params.Globals.Buildarch,
|
"arch", runtime.GOARCH,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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 (
|
||||||
@@ -37,8 +38,9 @@ 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 not yet written
|
||||||
// take the report files past the configured maximum size.
|
// 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")
|
var ErrFull = errors.New("report files at their size cap")
|
||||||
|
|
||||||
// Params defines the dependencies for Buffer.
|
// Params defines the dependencies for Buffer.
|
||||||
@@ -49,15 +51,28 @@ 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
|
// files are the report files that may be deleted to make room,
|
||||||
maxBytes int64
|
// oldest first: those in dataDir at start, then each one this
|
||||||
mu sync.Mutex
|
// 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
|
// 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
|
||||||
@@ -97,12 +112,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 +159,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 +174,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 +205,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 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
|
// 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() {
|
||||||
@@ -270,25 +335,28 @@ func (b *Buffer) writeFile(data []byte) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// The reports counted at their uncompressed size while they
|
// The reports counted at their uncompressed size while they
|
||||||
// waited; now they count as the file. After a failed write they
|
// waited; now they count as the file, which from here on may be
|
||||||
// stay counted as they were, which errs toward refusing reports
|
// deleted to make room. After a failed write they stay counted as
|
||||||
// early rather than letting the files pass the cap.
|
// they were, which errs toward refusing reports early rather than
|
||||||
|
// letting the files pass the cap.
|
||||||
b.mu.Lock()
|
b.mu.Lock()
|
||||||
b.usedBytes += info.Size() - int64(len(data))
|
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()
|
b.mu.Unlock()
|
||||||
|
|
||||||
return nil
|
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 +367,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
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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,124 @@ 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")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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
|
// TestConcurrentAppendsStopAtCap appends from many goroutines at once
|
||||||
@@ -407,6 +602,14 @@ func readReportFiles(t *testing.T, dir string) []string {
|
|||||||
return contents
|
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) {
|
func writeBytes(t *testing.T, path string, n int) {
|
||||||
t.Helper()
|
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 {
|
func hasReportFile(t *testing.T, dir string) bool {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"net"
|
"net"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
"runtime"
|
||||||
"strconv"
|
"strconv"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -53,7 +54,7 @@ func (s *Server) listenAndServe() {
|
|||||||
s.log.Info("http begin listen",
|
s.log.Info("http begin listen",
|
||||||
"listenaddr", s.httpServer.Addr,
|
"listenaddr", s.httpServer.Addr,
|
||||||
"version", s.params.Globals.Version,
|
"version", s.params.Globals.Version,
|
||||||
"buildarch", s.params.Globals.Buildarch,
|
"arch", runtime.GOARCH,
|
||||||
)
|
)
|
||||||
|
|
||||||
err := s.httpServer.ListenAndServe()
|
err := s.httpServer.ListenAndServe()
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
#!/bin/sh
|
#!/bin/sh
|
||||||
# script/build: compile the static netwatch-server binary into the
|
# script/build: compile the static netwatch-server binary into the
|
||||||
# backend project root, with its version and architecture stamped in.
|
# backend project root, with its version stamped in.
|
||||||
set -eu
|
set -eu
|
||||||
|
|
||||||
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
||||||
@@ -14,7 +14,7 @@ main() {
|
|||||||
version="${VERSION:-$(git describe --always --dirty 2>/dev/null || echo dev)}"
|
version="${VERSION:-$(git describe --always --dirty 2>/dev/null || echo dev)}"
|
||||||
|
|
||||||
CGO_ENABLED=0 go build -trimpath \
|
CGO_ENABLED=0 go build -trimpath \
|
||||||
-ldflags "-s -w -X main.Version=$version -X main.Buildarch=$(uname -m)" \
|
-ldflags "-s -w -X main.Version=$version" \
|
||||||
-o netwatch-server ./cmd/netwatch-server/
|
-o netwatch-server ./cmd/netwatch-server/
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+2
-3
@@ -16,9 +16,8 @@ main() {
|
|||||||
"$SCRIPT_DIR/check"
|
"$SCRIPT_DIR/check"
|
||||||
# Own line: a failing command substitution inside an argument does
|
# Own line: a failing command substitution inside an argument does
|
||||||
# not trip `set -e`, so the inline form degrades silently to an
|
# not trip `set -e`, so the inline form degrades silently to an
|
||||||
# empty constant. VERSION is computed here because .dockerignore
|
# empty constant. The VERSION build argument takes precedence over
|
||||||
# excludes .git, so `git describe` in a build stage yields an empty
|
# the version a build stage derives from the .git in the context.
|
||||||
# version without failing.
|
|
||||||
version="$(git describe --tags --always --dirty 2>/dev/null || true)"
|
version="$(git describe --tags --always --dirty 2>/dev/null || true)"
|
||||||
[ -n "$version" ] || version="unknown"
|
[ -n "$version" ] || version="unknown"
|
||||||
docker build --no-cache \
|
docker build --no-cache \
|
||||||
|
|||||||
+2
-3
@@ -12,9 +12,8 @@ main() {
|
|||||||
cd "$ROOT"
|
cd "$ROOT"
|
||||||
# Own line: a failing command substitution inside an argument does
|
# Own line: a failing command substitution inside an argument does
|
||||||
# not trip `set -e`, so the inline form degrades silently to an
|
# not trip `set -e`, so the inline form degrades silently to an
|
||||||
# empty constant. VERSION is computed here because .dockerignore
|
# empty constant. The VERSION build argument takes precedence over
|
||||||
# excludes .git, so `git describe` in a build stage yields an empty
|
# the version a build stage derives from the .git in the context.
|
||||||
# version without failing.
|
|
||||||
version="$(git describe --tags --always --dirty 2>/dev/null || true)"
|
version="$(git describe --tags --always --dirty 2>/dev/null || true)"
|
||||||
[ -n "$version" ] || version="unknown"
|
[ -n "$version" ] || version="unknown"
|
||||||
docker build --no-cache \
|
docker build --no-cache \
|
||||||
|
|||||||
Reference in New Issue
Block a user