next → main: data directory safe as root, version from git, 80% target timeout, report retention #85
@@ -4,3 +4,6 @@ tmp
|
||||
.DS_Store
|
||||
*.log
|
||||
.claude
|
||||
|
||||
# .git is sent so the build can stamp the version, without its config.
|
||||
.git/config
|
||||
|
||||
+23
-6
@@ -20,7 +20,10 @@ RUN make lint
|
||||
# golang:1.25-alpine (2026-02-27)
|
||||
FROM golang:1.25-alpine@sha256:f6751d823c26342f9506c03797d2527668d095b0a15f1862cddb4d927a7a4ced AS builder
|
||||
|
||||
RUN apk add --no-cache make
|
||||
# gcc and musl-dev are for make test: its race detector needs cgo, which
|
||||
# Go turns on by itself once a C compiler is present. make build still
|
||||
# sets CGO_ENABLED=0, so the binary stays static.
|
||||
RUN apk add --no-cache gcc git make musl-dev
|
||||
|
||||
WORKDIR /src
|
||||
|
||||
@@ -37,11 +40,24 @@ RUN make test
|
||||
|
||||
# make build is a shim around backend/script/build, the one definition
|
||||
# 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
|
||||
# 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
|
||||
# node:22-alpine as of 2026-02-22
|
||||
@@ -52,8 +68,9 @@ RUN yarn install --frozen-lockfile
|
||||
RUN apk add --no-cache git make
|
||||
COPY . .
|
||||
# make frontend-check is the frontend half of make check (test + lint +
|
||||
# fmt-check); its test step is the production yarn build, so this both
|
||||
# produces dist/ and gates the image on lint/fmt-check/test regressions.
|
||||
# fmt-check); its test step runs the unit tests, then the production
|
||||
# yarn build, so this both produces dist/ and gates the image on
|
||||
# lint/fmt-check/test regressions.
|
||||
# This node stage has neither Go nor Docker; the lint and builder stages
|
||||
# above gate the backend half.
|
||||
RUN make frontend-check
|
||||
|
||||
@@ -39,21 +39,23 @@ halves, so the root `make check` fails if either one is broken. We provide:
|
||||
- `script/bootstrap` — install all dependencies (the pinned node via nvm unless
|
||||
one new enough for the frontend's dependencies is installed, yarn via
|
||||
corepack, `yarn install --frozen-lockfile`, the pinned Go unless one at least
|
||||
as new as `backend/go.mod` asks for is installed, and the Go modules), linking
|
||||
what it installs itself into `~/.local/bin`, which has to be on `PATH`. It
|
||||
installs no Go linter and not Docker: `make lint` runs the linter in Docker
|
||||
as new as `backend/go.mod` asks for is installed, the Go modules, and gcc with
|
||||
the C library headers unless gcc is installed, for the race detector in
|
||||
`make test`), linking what it installs itself into `~/.local/bin`, which has
|
||||
to be on `PATH`. It installs no Go linter and not Docker: `make lint` runs the
|
||||
linter in Docker
|
||||
- `script/setup` — make a fresh clone ready for development: bootstrap plus the
|
||||
git pre-commit hook
|
||||
- `script/projectname` — print the project name (used for the Docker image tag)
|
||||
- `script/test` — run `script/frontend-test`, then the backend's Go tests, both
|
||||
within one 30-second timeout
|
||||
- `script/test` — run `script/frontend-test`, then the backend's Go tests, each
|
||||
under its own 30-second timeout
|
||||
- `script/lint` — run `script/frontend-lint`, then golangci-lint in Docker, by
|
||||
building the lint stage of `Dockerfile` without the cache
|
||||
- `script/fmt` — format all files (writes): prettier, then gofmt over `backend/`
|
||||
- `script/fmt-check` — check formatting (read-only): prettier, then gofmt
|
||||
- `script/check` — run test, lint, and fmt-check
|
||||
- `script/frontend-test` — run the production build as the frontend's test (no
|
||||
unit tests yet)
|
||||
- `script/frontend-test` — run the unit tests in `test/unit/` with Node's
|
||||
built-in test runner, then the production build
|
||||
- `script/frontend-lint` — run prettier in check mode
|
||||
- `script/frontend-fmt` — format everything prettier understands (writes)
|
||||
- `script/frontend-fmt-check` — check prettier formatting (read-only)
|
||||
@@ -136,8 +138,14 @@ Local hosts are tracked separately from WAN stats.
|
||||
### Latency measurement
|
||||
|
||||
HEAD requests with `mode: 'no-cors'` and `cache: 'no-store'`, timed with
|
||||
`performance.now()`. 1-second timeout; anything over 1000ms is clamped to
|
||||
unreachable. IPv4 only.
|
||||
`performance.now()`. Each check times out after 80% of the refresh interval (24
|
||||
seconds at 30 seconds) and is then recorded as a timeout, so a round's checks
|
||||
have all finished before the next round is due. When no WAN host answers, a
|
||||
recovery probe checks 4 random WAN hosts every half second, giving up the checks
|
||||
it started half a second before. As soon as one answers, a new round starts at
|
||||
once, as it does after an interval change. A round started early gives up the
|
||||
last round's checks if they are still waiting, and that round records nothing,
|
||||
so rounds never overlap. IPv4 only.
|
||||
|
||||
### Color coding
|
||||
|
||||
@@ -212,12 +220,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
|
||||
minute
|
||||
- `DATA_DIR_MAX_BYTES`, default `1073741824` (1 GiB): the most room the
|
||||
report files may take
|
||||
report files may take; the oldest are deleted to stay under it
|
||||
- `CORS_ALLOWED_ORIGINS`, default empty: other origins whose pages may call
|
||||
the API
|
||||
- `DEBUG`, default `false`: debug logging
|
||||
- `DATA_DIR`, default `/data/reports`: leave unset; reports kept outside
|
||||
`/data` do not survive a redeploy
|
||||
- `DATA_DIR`, default `/data/reports`: the directory the reports are kept
|
||||
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
|
||||
in front of the container connects from, as an IP address or CIDR; several
|
||||
are separated by commas. nginx takes the client address from
|
||||
|
||||
@@ -23,6 +23,52 @@ latest run passes.
|
||||
|
||||
# Completed Steps
|
||||
|
||||
- 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
|
||||
`.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
|
||||
are the same, the org standard changed, so update `GOLANGCI_CONFIG_SHA256` in
|
||||
that script. It used to say only to restore the file, which loops once the org
|
||||
standard itself has moved. A missing `.golangci.yml`, and a `sha256sum` that
|
||||
is missing or prints no hash, each get their own message instead of being
|
||||
reported as a mismatch; every one still fails the lint
|
||||
- 2026-10-03: the Go tests run with the race detector and coverage (issue #88):
|
||||
`backend/script/test` runs `go test -timeout 30s -race -cover ./...` and, if
|
||||
that fails, runs it again with `-v` and fails. Go's `-timeout` bounds the
|
||||
tests, not their compile; the root `script/test` no longer puts one 30-second
|
||||
timeout around both halves, which a cold Go build cache could use up on
|
||||
compiling alone. The race detector needs a C compiler: the builder stage of
|
||||
`Dockerfile` has gcc and musl-dev, and `script/bootstrap` installs gcc, with
|
||||
the C library headers on apt and apk, when gcc is missing; the binary is still
|
||||
built with `CGO_ENABLED=0`. New tests cover the health check's answer, a valid
|
||||
report's answer, a report file's exact contents, and the flush when the buffer
|
||||
reaches 10 MiB; the handlers' `TestImport` stub is gone
|
||||
- 2026-10-03: each target check times out after 80% of the refresh interval
|
||||
(issue #78), 24 seconds at 30 seconds, where it was capped at 3 seconds. A
|
||||
round started early, after an interval change or when the recovery probe finds
|
||||
a target answering, gives up the last round's checks if they are still
|
||||
waiting, so rounds never overlap; the recovery probe gives up its own checks
|
||||
after half a second. The frontend has its first unit tests, run by
|
||||
`script/frontend-test` with Node's built-in test runner; for them,
|
||||
`index.html` now links `src/styles.css`, which `src/main.js` used to import
|
||||
- 2026-09-29: the container sets up its own data directory (issue #75):
|
||||
`bin/entrypoint.sh`, still as root, creates `DATA_DIR` if missing and gives it
|
||||
and `/data` to the `netwatch` user with mode 750 before starting the backend
|
||||
|
||||
+25
-15
@@ -32,10 +32,13 @@ pattern as the repo root: the targets in `backend/Makefile` are thin shims over
|
||||
`test`, `fmt` and `fmt-check`:
|
||||
|
||||
- `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
|
||||
checkout, then to `dev`
|
||||
- `script/test` — run the Go tests under a 30-second timeout
|
||||
- `script/test` — run the Go tests with the race detector and coverage. Go's
|
||||
`-timeout 30s` bounds the tests, not their compile. If they fail, they run
|
||||
again with `-v` for the details, and the script fails. The race detector needs
|
||||
a C compiler
|
||||
- `script/lint` — check `.golangci.yml` against its pinned sha256, then run
|
||||
golangci-lint. It runs inside the golangci-lint image of the lint stage of
|
||||
the root `Dockerfile`; from a checkout, run `make lint` at the repo root,
|
||||
@@ -80,7 +83,7 @@ Internal packages in `internal/` follow standard Go project layout:
|
||||
| `BIND_ADDRESS` | empty | IP address to listen on; empty listens on every interface |
|
||||
| `PORT` | `8080` | HTTP listen port |
|
||||
| `DATA_DIR` | `./data/reports` | Directory for compressed reports |
|
||||
| `DATA_DIR_MAX_BYTES` | `1073741824` (1 GiB) | Largest total size of the report files in `DATA_DIR`; see [Report limits](#report-limits) |
|
||||
| `DATA_DIR_MAX_BYTES` | `1073741824` (1 GiB) | Most bytes of report files kept in `DATA_DIR`, oldest deleted first; see [Report limits](#report-limits) |
|
||||
| `DEBUG` | `false` | Enable debug logging |
|
||||
| `TRUSTED_PROXIES` | loopback + RFC1918 | Comma-separated CIDRs whose `X-Forwarded-For` / `X-Real-IP` headers are trusted for client IP resolution |
|
||||
| `REPORTS_PER_MINUTE` | `60` | Reports each client address may send a minute; see [Report limits](#report-limits) |
|
||||
@@ -104,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
|
||||
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
|
||||
the `/data` volume; the entrypoint creates it and gives it and `/data` to
|
||||
`netwatch` before starting the server. nginx replaces the security headers
|
||||
the `/data` volume; before starting the server, the entrypoint creates it and
|
||||
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
|
||||
what clients of the image see.
|
||||
|
||||
@@ -125,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
|
||||
time. The number starts at 1 when the server starts and goes up by one for each
|
||||
file the server starts to write, so two files written in the same millisecond
|
||||
still get different names. A failed write uses up its number, leaving a gap in
|
||||
the numbers if the file could not be created and otherwise a file under that
|
||||
number that may be incomplete. Each file contains one JSON object per line,
|
||||
compressed with zstd. Files are created with `O_EXCL` to prevent overwrites.
|
||||
still get different names. A failed write uses up its number and leaves a gap in
|
||||
the numbers: its file, if it was created, is removed. The file stays, counted
|
||||
toward `DATA_DIR_MAX_BYTES` from the next start, only if removing it fails too.
|
||||
Each file contains one JSON object per line, compressed with zstd. Files are
|
||||
created with `O_EXCL` to prevent overwrites.
|
||||
|
||||
### Report limits
|
||||
|
||||
@@ -148,11 +153,17 @@ credentials, so it is bounded instead. Both refusals below answer with the same
|
||||
`X-RateLimit-Reset` headers.
|
||||
- **Size cap.** The report files in `DATA_DIR` may total at most
|
||||
`DATA_DIR_MAX_BYTES`, counting the files already there at start. Reports
|
||||
waiting in memory count at their uncompressed size until they are written, so
|
||||
a report that would take the total past the cap is refused with 507, and
|
||||
nothing of it is stored. Deleting report files frees room only at the next
|
||||
start, when the files are counted again. The default of 1 GiB is small enough
|
||||
for any host; set it to the space you can give `DATA_DIR`.
|
||||
waiting in memory count at their uncompressed size until they are written;
|
||||
those lost to a failed write stop counting, and the part of its file written
|
||||
is removed. When a report would take the total past the cap, the oldest report
|
||||
files are deleted to make room, and each deletion is logged with the file's
|
||||
name and size; a file still being written is never deleted. A report is
|
||||
refused with 507, and nothing of it is stored, only when the reports waiting
|
||||
to be written fill the cap on their own, and then no file is deleted. At
|
||||
start, report files past the cap, as after lowering it, are deleted the same
|
||||
way. So the cap is how much of the newest reports is kept: the default of 1
|
||||
GiB is small enough for any host; set it to the space you can give
|
||||
`DATA_DIR`.
|
||||
|
||||
### CORS
|
||||
|
||||
@@ -170,7 +181,6 @@ starting, with an error naming `CORS_ALLOWED_ORIGINS`.
|
||||
- Add integration test that POSTs a report and verifies the compressed output
|
||||
- Add report decompression/query endpoint
|
||||
- Add metrics (Prometheus) for buffer size, flush count, report count
|
||||
- Add retention policy to prune old report files
|
||||
|
||||
## License
|
||||
|
||||
|
||||
@@ -4,6 +4,7 @@ package main
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"os/user"
|
||||
|
||||
"sneak.berlin/go/netwatch/internal/config"
|
||||
"sneak.berlin/go/netwatch/internal/globals"
|
||||
@@ -19,9 +20,8 @@ import (
|
||||
|
||||
//nolint:gochecknoglobals // set via ldflags at build time
|
||||
var (
|
||||
Appname = "netwatch-server"
|
||||
Version string
|
||||
Buildarch string
|
||||
Appname = "netwatch-server"
|
||||
Version string
|
||||
)
|
||||
|
||||
func main() {
|
||||
@@ -38,9 +38,28 @@ func main() {
|
||||
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.Version = Version
|
||||
globals.Buildarch = Buildarch
|
||||
|
||||
fx.New(
|
||||
fx.Provide(
|
||||
|
||||
@@ -10,22 +10,18 @@ var (
|
||||
Appname string
|
||||
// Version is the git version tag.
|
||||
Version string
|
||||
// Buildarch is the build architecture.
|
||||
Buildarch string
|
||||
)
|
||||
|
||||
// Globals holds build-time metadata for the application.
|
||||
type Globals struct {
|
||||
Appname string
|
||||
Version string
|
||||
Buildarch string
|
||||
Appname string
|
||||
Version string
|
||||
}
|
||||
|
||||
// New creates a Globals instance from package-level variables.
|
||||
func New(_ fx.Lifecycle) (*Globals, error) {
|
||||
return &Globals{
|
||||
Appname: Appname,
|
||||
Buildarch: Buildarch,
|
||||
Version: Version,
|
||||
Appname: Appname,
|
||||
Version: Version,
|
||||
}, nil
|
||||
}
|
||||
|
||||
@@ -1,13 +0,0 @@
|
||||
package handlers_test
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
_ "sneak.berlin/go/netwatch/internal/handlers"
|
||||
)
|
||||
|
||||
func TestImport(t *testing.T) {
|
||||
t.Parallel()
|
||||
// Compilation check — verifies the package parses
|
||||
// and all imports resolve.
|
||||
}
|
||||
@@ -0,0 +1,120 @@
|
||||
package handlers_test
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"maps"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"slices"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"sneak.berlin/go/netwatch/internal/globals"
|
||||
"sneak.berlin/go/netwatch/internal/handlers"
|
||||
"sneak.berlin/go/netwatch/internal/healthcheck"
|
||||
"sneak.berlin/go/netwatch/internal/logger"
|
||||
|
||||
"go.uber.org/fx/fxtest"
|
||||
)
|
||||
|
||||
// newStartedHandlers builds Handlers with a real health check for the
|
||||
// server named in g, and starts them, which records the time the
|
||||
// uptime counts from.
|
||||
func newStartedHandlers(t *testing.T, g *globals.Globals) *handlers.Handlers {
|
||||
t.Helper()
|
||||
|
||||
lc := fxtest.NewLifecycle(t)
|
||||
|
||||
log, err := logger.New(lc, logger.Params{Globals: g})
|
||||
if err != nil {
|
||||
t.Fatalf("logger: %v", err)
|
||||
}
|
||||
|
||||
hc, err := healthcheck.New(lc,
|
||||
healthcheck.Params{Globals: g, Logger: log})
|
||||
if err != nil {
|
||||
t.Fatalf("health check: %v", err)
|
||||
}
|
||||
|
||||
h, err := handlers.New(lc,
|
||||
handlers.Params{Globals: g, Healthcheck: hc, Logger: log})
|
||||
if err != nil {
|
||||
t.Fatalf("handlers: %v", err)
|
||||
}
|
||||
|
||||
lc.RequireStart()
|
||||
t.Cleanup(lc.RequireStop)
|
||||
|
||||
return h
|
||||
}
|
||||
|
||||
// TestHandleHealthCheck checks the health check's answer: 200, a JSON
|
||||
// content type, and a JSON object with exactly the fields of
|
||||
// healthcheck.Response, carrying this server's name and version and
|
||||
// an uptime counted from its start.
|
||||
func TestHandleHealthCheck(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
g := &globals.Globals{Appname: "netwatch-server", Version: "v1.2.3"}
|
||||
h := newStartedHandlers(t, g)
|
||||
|
||||
rec := httptest.NewRecorder()
|
||||
req := httptest.NewRequestWithContext(t.Context(),
|
||||
http.MethodGet, "/.well-known/healthcheck", http.NoBody)
|
||||
|
||||
h.HandleHealthCheck().ServeHTTP(rec, req)
|
||||
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("status = %d, want %d", rec.Code, http.StatusOK)
|
||||
}
|
||||
|
||||
contentType := rec.Header().Get("Content-Type")
|
||||
if contentType != "application/json; charset=utf-8" {
|
||||
t.Errorf("Content-Type = %q, want %q",
|
||||
contentType, "application/json; charset=utf-8")
|
||||
}
|
||||
|
||||
var body map[string]any
|
||||
|
||||
err := json.Unmarshal(rec.Body.Bytes(), &body)
|
||||
if err != nil {
|
||||
t.Fatalf("body not a JSON object: %v (%q)", err, rec.Body.String())
|
||||
}
|
||||
|
||||
fields := []string{
|
||||
"appname", "now", "status", "uptimeHuman", "uptimeSeconds", "version",
|
||||
}
|
||||
if got := slices.Sorted(maps.Keys(body)); !slices.Equal(got, fields) {
|
||||
t.Fatalf("fields = %v, want %v", got, fields)
|
||||
}
|
||||
|
||||
for field, want := range map[string]string{
|
||||
"appname": g.Appname, "status": "ok", "version": g.Version,
|
||||
} {
|
||||
if body[field] != want {
|
||||
t.Errorf("%s = %v, want %q", field, body[field], want)
|
||||
}
|
||||
}
|
||||
|
||||
now, _ := body["now"].(string)
|
||||
|
||||
at, err := time.Parse(time.RFC3339Nano, now)
|
||||
if err != nil || time.Since(at).Abs() > time.Minute {
|
||||
t.Errorf("now = %q, want the current time in RFC 3339 (%v)", now, err)
|
||||
}
|
||||
|
||||
// Started just now, so the uptime is well under a minute.
|
||||
human, _ := body["uptimeHuman"].(string)
|
||||
|
||||
uptime, err := time.ParseDuration(human)
|
||||
if err != nil || uptime > time.Minute {
|
||||
t.Errorf("uptimeHuman = %q, want a duration under a minute (%v)",
|
||||
human, err)
|
||||
}
|
||||
|
||||
seconds, ok := body["uptimeSeconds"].(float64)
|
||||
if !ok || seconds < 0 || seconds > time.Minute.Seconds() {
|
||||
t.Errorf("uptimeSeconds = %v, want a number of seconds under a minute",
|
||||
body["uptimeSeconds"])
|
||||
}
|
||||
}
|
||||
@@ -86,11 +86,12 @@ func (s *Handlers) decodeErrorStatus(err error) int {
|
||||
}
|
||||
|
||||
// appendErrorStatus logs a failure to store a report and returns
|
||||
// the status to send: 507 when the report files are at their size
|
||||
// cap, otherwise 500.
|
||||
// the status to send: 507 when the reports waiting to be written fill
|
||||
// the size cap, otherwise 500.
|
||||
func (s *Handlers) appendErrorStatus(err error) int {
|
||||
if errors.Is(err, reportbuf.ErrFull) {
|
||||
s.log.Warn("report refused: report files at their size cap")
|
||||
s.log.Warn("report refused: " +
|
||||
"reports waiting to be written fill the size cap")
|
||||
|
||||
return http.StatusInsufficientStorage
|
||||
}
|
||||
|
||||
@@ -46,6 +46,63 @@ func decodeStatus(t *testing.T, body []byte) string {
|
||||
return resp.Status
|
||||
}
|
||||
|
||||
// TestHandleReportAcceptsValidReports checks the answer to a valid
|
||||
// report, one with no hosts and one shaped as the frontend sends them:
|
||||
// 200 and {"status":"ok"} as JSON.
|
||||
func TestHandleReportAcceptsValidReports(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
body string
|
||||
}{
|
||||
{
|
||||
name: "no hosts",
|
||||
body: `{"clientId":"c1","geo":null,"hosts":[],` +
|
||||
`"timestamp":"2026-10-03T12:00:00.000Z"}`,
|
||||
},
|
||||
{
|
||||
name: "a host with a latency and an error sample",
|
||||
body: `{"clientId":"c1","geo":null,"hosts":[{` +
|
||||
`"name":"Example","url":"https://example.com/",` +
|
||||
`"status":"error","history":[` +
|
||||
`{"t":1790000000000,"latency":42,"error":null},` +
|
||||
`{"t":1790000003000,"latency":null,"error":"timeout"}]}],` +
|
||||
`"timestamp":"2026-10-03T12:00:00.000Z"}`,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
h := newTestHandlers(stubAppender{}, io.Discard)
|
||||
|
||||
rec := httptest.NewRecorder()
|
||||
req := httptest.NewRequestWithContext(t.Context(),
|
||||
http.MethodPost, "/api/v1/reports",
|
||||
strings.NewReader(tt.body),
|
||||
)
|
||||
|
||||
h.HandleReport().ServeHTTP(rec, req)
|
||||
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("status = %d, want %d", rec.Code, http.StatusOK)
|
||||
}
|
||||
|
||||
contentType := rec.Header().Get("Content-Type")
|
||||
if contentType != "application/json; charset=utf-8" {
|
||||
t.Errorf("Content-Type = %q, want %q",
|
||||
contentType, "application/json; charset=utf-8")
|
||||
}
|
||||
|
||||
if got := rec.Body.String(); got != "{\"status\":\"ok\"}\n" {
|
||||
t.Errorf("body = %q, want %q", got, "{\"status\":\"ok\"}\n")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestHandleReportStorageFailureIsNon2xx(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
@@ -68,9 +125,9 @@ func TestHandleReportStorageFailureIsNon2xx(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestHandleReportFullIs507 checks the answer when the report files
|
||||
// are at their size cap: 507 and the usual error body, which tells
|
||||
// the client nothing more.
|
||||
// TestHandleReportFullIs507 checks the answer when the reports waiting
|
||||
// to be written fill the size cap: 507 and the usual error body, which
|
||||
// tells the client nothing more.
|
||||
func TestHandleReportFullIs507(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
|
||||
@@ -5,6 +5,7 @@ package logger
|
||||
import (
|
||||
"log/slog"
|
||||
"os"
|
||||
"runtime"
|
||||
|
||||
"sneak.berlin/go/netwatch/internal/globals"
|
||||
|
||||
@@ -95,6 +96,6 @@ func (l *Logger) Identify() {
|
||||
l.log.Info("starting",
|
||||
"appname", l.params.Globals.Appname,
|
||||
"version", l.params.Globals.Version,
|
||||
"buildarch", l.params.Globals.Buildarch,
|
||||
"arch", runtime.GOARCH,
|
||||
)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -1,6 +1,13 @@
|
||||
package reportbuf
|
||||
|
||||
import "time"
|
||||
import (
|
||||
"os"
|
||||
"time"
|
||||
)
|
||||
|
||||
// FlushSizeThreshold exposes the buffer size at which Append starts
|
||||
// writing a report file to the external tests.
|
||||
const FlushSizeThreshold = flushSizeThreshold
|
||||
|
||||
// Flush writes the buffered reports to a file now, as the periodic
|
||||
// flush does, so tests need not wait a minute for it.
|
||||
@@ -13,3 +20,10 @@ func (b *Buffer) Flush() error {
|
||||
func (b *Buffer) StopClock(at time.Time) {
|
||||
b.now = func() time.Time { return at }
|
||||
}
|
||||
|
||||
// OnFileCreated makes the buffer call fn with each report file it
|
||||
// writes from now on, once the file is created and before anything is
|
||||
// written to it.
|
||||
func (b *Buffer) OnFileCreated(fn func(f *os.File)) {
|
||||
b.fileCreated = fn
|
||||
}
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
// Package reportbuf accumulates telemetry reports in memory
|
||||
// and periodically flushes them to zstd-compressed JSONL files.
|
||||
// and periodically flushes them to zstd-compressed JSONL files,
|
||||
// deleting the oldest files to keep them under a size cap.
|
||||
package reportbuf
|
||||
|
||||
import (
|
||||
@@ -12,6 +13,7 @@ import (
|
||||
"log/slog"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"slices"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
@@ -37,9 +39,10 @@ const (
|
||||
fileSuffix = ".jsonl.zst"
|
||||
)
|
||||
|
||||
// ErrFull is returned by Append when storing the report would
|
||||
// take the report files past the configured maximum size.
|
||||
var ErrFull = errors.New("report files at their size cap")
|
||||
// ErrFull is returned by Append when the reports waiting to be
|
||||
// written leave no room for the report under the configured maximum
|
||||
// size, however many report files are deleted.
|
||||
var ErrFull = errors.New("reports waiting to be written fill the size cap")
|
||||
|
||||
// Params defines the dependencies for Buffer.
|
||||
type Params struct {
|
||||
@@ -49,15 +52,34 @@ type Params struct {
|
||||
Logger *logger.Logger
|
||||
}
|
||||
|
||||
// reportFile is a report file that may be deleted to make room, with
|
||||
// the size it counts for in usedBytes.
|
||||
type reportFile struct {
|
||||
name string
|
||||
size int64
|
||||
}
|
||||
|
||||
// Buffer accumulates JSON lines in memory and flushes them
|
||||
// to zstd-compressed files on disk.
|
||||
type Buffer struct {
|
||||
buf bytes.Buffer
|
||||
dataDir string
|
||||
done chan struct{}
|
||||
log *slog.Logger
|
||||
maxBytes int64
|
||||
mu sync.Mutex
|
||||
buf bytes.Buffer
|
||||
dataDir string
|
||||
done chan struct{}
|
||||
// fileCreated is called with each report file once it is
|
||||
// created, before anything is written to it: it does nothing,
|
||||
// except in tests that hold the write open or make it fail.
|
||||
fileCreated func(f *os.File)
|
||||
// files are the report files that may be deleted to make room,
|
||||
// in name order, which is oldest first: those in dataDir at
|
||||
// start, and each one this buffer writes, put in at its place by
|
||||
// name once it is complete, even when an older file's write
|
||||
// completes after a newer one's. A file still being written is
|
||||
// not among them. filesBytes is their total size.
|
||||
files []reportFile
|
||||
filesBytes int64
|
||||
log *slog.Logger
|
||||
maxBytes int64
|
||||
mu sync.Mutex
|
||||
// now is the clock report files are named by: time.Now, except
|
||||
// in tests that need two flushes to share a timestamp.
|
||||
now func() time.Time
|
||||
@@ -66,8 +88,8 @@ type Buffer struct {
|
||||
seq atomic.Uint64
|
||||
stopOnce sync.Once
|
||||
// usedBytes is what Append checks against maxBytes: the size
|
||||
// of the report files in dataDir, plus the reports not yet
|
||||
// written to one at their uncompressed size.
|
||||
// of the report files in dataDir, plus the reports waiting to
|
||||
// be written to one at their uncompressed size.
|
||||
usedBytes int64
|
||||
}
|
||||
|
||||
@@ -83,11 +105,12 @@ func New(
|
||||
}
|
||||
|
||||
b := &Buffer{
|
||||
dataDir: dir,
|
||||
done: make(chan struct{}),
|
||||
log: params.Logger.Get(),
|
||||
maxBytes: params.Config.DataDirMaxBytes,
|
||||
now: time.Now,
|
||||
dataDir: dir,
|
||||
done: make(chan struct{}),
|
||||
fileCreated: func(*os.File) {},
|
||||
log: params.Logger.Get(),
|
||||
maxBytes: params.Config.DataDirMaxBytes,
|
||||
now: time.Now,
|
||||
}
|
||||
|
||||
lc.Append(fx.Hook{
|
||||
@@ -97,12 +120,28 @@ func New(
|
||||
return fmt.Errorf("create data dir: %w", err)
|
||||
}
|
||||
|
||||
// Report files left by earlier runs count too.
|
||||
b.usedBytes, err = reportFilesSize(b.dataDir)
|
||||
// Report files left by earlier runs count too, and are
|
||||
// the first to be deleted to make room.
|
||||
files, err := reportFiles(b.dataDir)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
b.mu.Lock()
|
||||
|
||||
b.files = files
|
||||
for _, f := range files {
|
||||
b.filesBytes += f.size
|
||||
}
|
||||
|
||||
b.usedBytes = b.filesBytes
|
||||
|
||||
// The files may be past the cap, if it was lowered since
|
||||
// the last run.
|
||||
b.deleteOldestFiles(0)
|
||||
|
||||
b.mu.Unlock()
|
||||
|
||||
go b.flushLoop()
|
||||
|
||||
return nil
|
||||
@@ -128,9 +167,10 @@ func New(
|
||||
}
|
||||
|
||||
// Append marshals v as a single JSON line and appends it to
|
||||
// the buffer. It stores nothing and returns ErrFull if the line
|
||||
// would take usedBytes past maxBytes. If the buffer reaches the
|
||||
// size threshold, it is drained and written to disk
|
||||
// the buffer. If the line would take usedBytes past maxBytes, the
|
||||
// oldest report files are deleted to make room; it stores nothing
|
||||
// and returns ErrFull if that cannot make room. If the buffer
|
||||
// reaches the size threshold, it is drained and written to disk
|
||||
// asynchronously.
|
||||
func (b *Buffer) Append(v any) error {
|
||||
line, err := json.Marshal(v)
|
||||
@@ -142,6 +182,8 @@ func (b *Buffer) Append(v any) error {
|
||||
|
||||
b.mu.Lock()
|
||||
|
||||
b.deleteOldestFiles(lineBytes)
|
||||
|
||||
if b.usedBytes+lineBytes > b.maxBytes {
|
||||
b.mu.Unlock()
|
||||
|
||||
@@ -171,6 +213,37 @@ func (b *Buffer) Append(v any) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// deleteOldestFiles deletes report files, oldest first, until n more
|
||||
// bytes fit under maxBytes. It deletes none when the reports waiting
|
||||
// to be written leave no room for n even with every file gone, since
|
||||
// that would lose the files for nothing. The caller must hold b.mu.
|
||||
func (b *Buffer) deleteOldestFiles(n int64) {
|
||||
for b.usedBytes+n > b.maxBytes && len(b.files) > 0 {
|
||||
if b.usedBytes-b.filesBytes+n > b.maxBytes {
|
||||
return
|
||||
}
|
||||
|
||||
f := b.files[0]
|
||||
b.files = b.files[1:]
|
||||
b.filesBytes -= f.size
|
||||
|
||||
// A file already gone, deleted by hand, has freed its room too.
|
||||
err := os.Remove(filepath.Join(b.dataDir, f.name))
|
||||
if err != nil && !errors.Is(err, fs.ErrNotExist) {
|
||||
// The file is still there, so it still counts. It is
|
||||
// not tried again until the next start.
|
||||
b.log.Error("delete report file failed",
|
||||
"file", f.name, "error", err)
|
||||
|
||||
continue
|
||||
}
|
||||
|
||||
b.usedBytes -= f.size
|
||||
b.log.Info("deleted report file to make room",
|
||||
"file", f.name, "bytes", f.size)
|
||||
}
|
||||
}
|
||||
|
||||
// flushLoop runs a ticker that periodically flushes buffered
|
||||
// data to disk until the done channel is closed.
|
||||
func (b *Buffer) flushLoop() {
|
||||
@@ -217,15 +290,48 @@ func (b *Buffer) drainBuf() []byte {
|
||||
return data
|
||||
}
|
||||
|
||||
// writeFile creates a timestamped zstd-compressed JSONL file
|
||||
// in the data directory.
|
||||
// writeFile writes data, reports drained from the buffer, to a new
|
||||
// timestamped zstd-compressed JSONL file in the data directory.
|
||||
func (b *Buffer) writeFile(data []byte) error {
|
||||
// The timestamp comes first, so the names sort by time; the number
|
||||
// after it tells apart files named in the same millisecond.
|
||||
ts := b.now().UTC().Format("2006-01-02T15-04-05.000Z")
|
||||
name := fmt.Sprintf("%s%s-%d%s", filePrefix, ts, b.seq.Add(1), fileSuffix)
|
||||
path := filepath.Join(b.dataDir, name)
|
||||
|
||||
size, err := b.createFile(filepath.Join(b.dataDir, name), data)
|
||||
|
||||
// The reports no longer wait to be written, so they stop counting
|
||||
// at their uncompressed size. If the write failed they are lost;
|
||||
// otherwise they count as the file, which from here on may be
|
||||
// deleted to make room.
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
|
||||
b.usedBytes -= int64(len(data))
|
||||
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
b.usedBytes += size
|
||||
b.filesBytes += size
|
||||
|
||||
// At its place by name, not at the end: another write, of a newer
|
||||
// file, may have completed while this one was being written.
|
||||
i, _ := slices.BinarySearchFunc(b.files, name,
|
||||
func(f reportFile, target string) int {
|
||||
return strings.Compare(f.name, target)
|
||||
})
|
||||
b.files = slices.Insert(b.files, i, reportFile{name: name, size: size})
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// createFile creates the file at path holding data compressed with
|
||||
// zstd, and returns its size. If the write fails once the file is
|
||||
// created, it removes the file, so that a failed write leaves nothing
|
||||
// behind to take room.
|
||||
func (b *Buffer) createFile(path string, data []byte) (int64, error) {
|
||||
// path is built from the operator-supplied dataDir plus a
|
||||
// generated timestamp and number, so it carries no external input.
|
||||
f, err := os.OpenFile( //nolint:gosec // see comment above
|
||||
@@ -234,61 +340,65 @@ func (b *Buffer) writeFile(data []byte) error {
|
||||
filePerms,
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("create report file: %w", err)
|
||||
return 0, fmt.Errorf("create report file: %w", err)
|
||||
}
|
||||
|
||||
// Closes the file on the early returns below. The success
|
||||
// path closes it explicitly to check the error; closing it
|
||||
// a second time here is harmless.
|
||||
defer func() { _ = f.Close() }()
|
||||
b.fileCreated(f)
|
||||
|
||||
size, err := writeCompressed(f, data)
|
||||
if err != nil {
|
||||
_ = f.Close()
|
||||
|
||||
return 0, errors.Join(err, os.Remove(path))
|
||||
}
|
||||
|
||||
err = f.Close()
|
||||
if err != nil {
|
||||
err = fmt.Errorf("close report file: %w", err)
|
||||
|
||||
return 0, errors.Join(err, os.Remove(path))
|
||||
}
|
||||
|
||||
return size, nil
|
||||
}
|
||||
|
||||
// writeCompressed writes data to f compressed with zstd, and returns
|
||||
// the size of f.
|
||||
func writeCompressed(f *os.File, data []byte) (int64, error) {
|
||||
enc, err := zstd.NewWriter(f)
|
||||
if err != nil {
|
||||
return fmt.Errorf("create zstd encoder: %w", err)
|
||||
return 0, fmt.Errorf("create zstd encoder: %w", err)
|
||||
}
|
||||
|
||||
_, err = enc.Write(data)
|
||||
if err != nil {
|
||||
_ = enc.Close()
|
||||
|
||||
return fmt.Errorf("write compressed data: %w", err)
|
||||
return 0, fmt.Errorf("write compressed data: %w", err)
|
||||
}
|
||||
|
||||
err = enc.Close()
|
||||
if err != nil {
|
||||
return fmt.Errorf("close zstd encoder: %w", err)
|
||||
return 0, fmt.Errorf("close zstd encoder: %w", err)
|
||||
}
|
||||
|
||||
info, err := f.Stat()
|
||||
if err != nil {
|
||||
return fmt.Errorf("stat report file: %w", err)
|
||||
return 0, fmt.Errorf("stat report file: %w", err)
|
||||
}
|
||||
|
||||
err = f.Close()
|
||||
if err != nil {
|
||||
return fmt.Errorf("close report file: %w", err)
|
||||
}
|
||||
|
||||
// The reports counted at their uncompressed size while they
|
||||
// waited; now they count as the file. After a failed write they
|
||||
// stay counted as they were, which errs toward refusing reports
|
||||
// early rather than letting the files pass the cap.
|
||||
b.mu.Lock()
|
||||
b.usedBytes += info.Size() - int64(len(data))
|
||||
b.mu.Unlock()
|
||||
|
||||
return nil
|
||||
return info.Size(), nil
|
||||
}
|
||||
|
||||
// reportFilesSize returns the total size of the report files in
|
||||
// dir.
|
||||
func reportFilesSize(dir string) (int64, error) {
|
||||
// reportFiles returns the report files in dir, oldest first:
|
||||
// os.ReadDir sorts them by name, and the names sort by time.
|
||||
func reportFiles(dir string) ([]reportFile, error) {
|
||||
entries, err := os.ReadDir(dir)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("read data dir: %w", err)
|
||||
return nil, fmt.Errorf("read data dir: %w", err)
|
||||
}
|
||||
|
||||
var total int64
|
||||
files := make([]reportFile, 0, len(entries))
|
||||
|
||||
for _, entry := range entries {
|
||||
name := entry.Name()
|
||||
@@ -299,11 +409,11 @@ func reportFilesSize(dir string) (int64, error) {
|
||||
|
||||
info, err := entry.Info()
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("stat report file: %w", err)
|
||||
return nil, fmt.Errorf("stat report file: %w", err)
|
||||
}
|
||||
|
||||
total += info.Size()
|
||||
files = append(files, reportFile{name: name, size: info.Size()})
|
||||
}
|
||||
|
||||
return total, nil
|
||||
return files, nil
|
||||
}
|
||||
|
||||
@@ -139,11 +139,19 @@ func lineBytes(t *testing.T, report any) int {
|
||||
return len(line) + 1
|
||||
}
|
||||
|
||||
// TestAppendPastCapIsRefused fills the cap with a report not yet
|
||||
// written. The next report is refused, and the report file already in
|
||||
// DATA_DIR is kept: it is smaller than a report, so deleting it could
|
||||
// not make room.
|
||||
func TestAppendPastCapIsRefused(t *testing.T) {
|
||||
report := map[string]string{"id": "cap"}
|
||||
dir := t.TempDir()
|
||||
earlier := reportFilePath(dir, 1)
|
||||
|
||||
t.Setenv("DATA_DIR", t.TempDir())
|
||||
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(lineBytes(t, report)))
|
||||
writeBytes(t, earlier, 1)
|
||||
|
||||
t.Setenv("DATA_DIR", dir)
|
||||
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(1+lineBytes(t, report)))
|
||||
|
||||
buf := startBuffer(t)
|
||||
|
||||
@@ -156,24 +164,38 @@ func TestAppendPastCapIsRefused(t *testing.T) {
|
||||
if !errors.Is(err, reportbuf.ErrFull) {
|
||||
t.Fatalf("report past the cap: error = %v, want ErrFull", err)
|
||||
}
|
||||
|
||||
if !exists(t, earlier) {
|
||||
t.Fatal("report file deleted, though that could not make room")
|
||||
}
|
||||
}
|
||||
|
||||
// TestCapCountsReportFilesAlreadyInDataDir starts on a data
|
||||
// directory holding a report file from an earlier run, and a file
|
||||
// that is not a report, which must not count.
|
||||
func TestCapCountsReportFilesAlreadyInDataDir(t *testing.T) {
|
||||
const earlierBytes = 100
|
||||
// TestOldestReportFileDeletedFirst starts on a data directory holding
|
||||
// report files from an earlier run, and a file that is not a report,
|
||||
// which neither counts nor is ever deleted. Nothing is deleted while
|
||||
// there is room; then only the oldest report file is.
|
||||
func TestOldestReportFileDeletedFirst(t *testing.T) {
|
||||
const fileBytes = 100
|
||||
|
||||
report := map[string]string{"id": "cap"}
|
||||
report := map[string]string{"id": "oldest"}
|
||||
dir := t.TempDir()
|
||||
oldest := reportFilePath(dir, 1)
|
||||
kept := []string{
|
||||
reportFilePath(dir, 2),
|
||||
reportFilePath(dir, 3),
|
||||
filepath.Join(dir, "notes.txt"),
|
||||
}
|
||||
|
||||
writeBytes(t, filepath.Join(dir, "reports-2026-01-01T00-00-00.000Z.jsonl.zst"),
|
||||
earlierBytes)
|
||||
writeBytes(t, filepath.Join(dir, "notes.txt"), 10*earlierBytes)
|
||||
writeBytes(t, oldest, fileBytes)
|
||||
|
||||
for _, path := range kept {
|
||||
writeBytes(t, path, fileBytes)
|
||||
}
|
||||
|
||||
t.Setenv("DATA_DIR", dir)
|
||||
// Room for the three report files and one report.
|
||||
t.Setenv("DATA_DIR_MAX_BYTES",
|
||||
strconv.Itoa(earlierBytes+lineBytes(t, report)))
|
||||
strconv.Itoa(3*fileBytes+lineBytes(t, report)))
|
||||
|
||||
buf := startBuffer(t)
|
||||
|
||||
@@ -182,9 +204,55 @@ func TestCapCountsReportFilesAlreadyInDataDir(t *testing.T) {
|
||||
t.Fatalf("report that fills the cap exactly: %v", err)
|
||||
}
|
||||
|
||||
if !exists(t, oldest) {
|
||||
t.Fatal("oldest report file deleted while there was room")
|
||||
}
|
||||
|
||||
err = buf.Append(report)
|
||||
if !errors.Is(err, reportbuf.ErrFull) {
|
||||
t.Fatalf("report past the cap: error = %v, want ErrFull", err)
|
||||
if err != nil {
|
||||
t.Fatalf("report past the cap: %v", err)
|
||||
}
|
||||
|
||||
if exists(t, oldest) {
|
||||
t.Fatal("oldest report file kept when room was needed")
|
||||
}
|
||||
|
||||
for _, path := range kept {
|
||||
if !exists(t, path) {
|
||||
t.Fatalf("%s deleted; only the oldest report file should be", path)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestStartDeletesFilesPastCap starts on report files past the cap,
|
||||
// as after the cap is lowered: the oldest are deleted until the rest
|
||||
// fit.
|
||||
func TestStartDeletesFilesPastCap(t *testing.T) {
|
||||
const fileBytes = 100
|
||||
|
||||
dir := t.TempDir()
|
||||
oldest := reportFilePath(dir, 1)
|
||||
kept := []string{reportFilePath(dir, 2), reportFilePath(dir, 3)}
|
||||
|
||||
writeBytes(t, oldest, fileBytes)
|
||||
|
||||
for _, path := range kept {
|
||||
writeBytes(t, path, fileBytes)
|
||||
}
|
||||
|
||||
t.Setenv("DATA_DIR", dir)
|
||||
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(len(kept)*fileBytes))
|
||||
|
||||
startBuffer(t)
|
||||
|
||||
if exists(t, oldest) {
|
||||
t.Fatal("oldest report file kept, though the files were past the cap")
|
||||
}
|
||||
|
||||
for _, path := range kept {
|
||||
if !exists(t, path) {
|
||||
t.Fatalf("%s deleted, though the rest fit without it", path)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -195,8 +263,9 @@ func TestWrittenReportsCountAtFileSize(t *testing.T) {
|
||||
// Repetitive, so its file is far smaller than its JSON.
|
||||
report := map[string]string{"id": strings.Repeat("a", 1000)}
|
||||
size := lineBytes(t, report)
|
||||
dir := t.TempDir()
|
||||
|
||||
t.Setenv("DATA_DIR", t.TempDir())
|
||||
t.Setenv("DATA_DIR", dir)
|
||||
// Room for the report twice over only if the first one counts
|
||||
// at its file's size by the time the second arrives.
|
||||
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(2*size-1))
|
||||
@@ -217,13 +286,22 @@ func TestWrittenReportsCountAtFileSize(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatalf("second report, after the first was written: %v", err)
|
||||
}
|
||||
|
||||
if !hasReportFile(t, dir) {
|
||||
t.Fatal("first report's file deleted, though the second fit beside it")
|
||||
}
|
||||
}
|
||||
|
||||
// TestWrittenReportsKeepCounting writes one report file after another
|
||||
// under a small cap: each report must be taken while the files on disk
|
||||
// leave room for it, and refused once they do not.
|
||||
func TestWrittenReportsKeepCounting(t *testing.T) {
|
||||
const maxBytes = 200
|
||||
// TestWrittenFilesDeletedToMakeRoom writes one report file after
|
||||
// another under a small cap. Every report must be taken; files are
|
||||
// deleted only when the report would not fit beside them, and the
|
||||
// files kept leave room for it.
|
||||
func TestWrittenFilesDeletedToMakeRoom(t *testing.T) {
|
||||
const (
|
||||
maxBytes = 200
|
||||
// One file each, which take far more than maxBytes together.
|
||||
reports = 50
|
||||
)
|
||||
|
||||
report := map[string]string{"id": "written"}
|
||||
size := int64(lineBytes(t, report))
|
||||
@@ -234,23 +312,24 @@ func TestWrittenReportsKeepCounting(t *testing.T) {
|
||||
|
||||
buf := startBuffer(t)
|
||||
|
||||
// Every file takes at least a byte, so they fill the cap within
|
||||
// maxBytes rounds.
|
||||
for range maxBytes {
|
||||
used := reportFilesBytes(t, dir)
|
||||
for range reports {
|
||||
before := reportFilesBytes(t, dir)
|
||||
|
||||
err := buf.Append(report)
|
||||
if used+size > maxBytes {
|
||||
if !errors.Is(err, reportbuf.ErrFull) {
|
||||
t.Fatalf("with %d bytes of report files: error = %v, "+
|
||||
"want ErrFull", used, err)
|
||||
}
|
||||
|
||||
return
|
||||
if err != nil {
|
||||
t.Fatalf("with %d bytes of report files: %v", before, err)
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
t.Fatalf("with %d bytes of report files: %v", used, err)
|
||||
after := reportFilesBytes(t, dir)
|
||||
|
||||
if before+size <= maxBytes && after != before {
|
||||
t.Fatalf("files deleted, though the report fit beside "+
|
||||
"their %d bytes", before)
|
||||
}
|
||||
|
||||
if after+size > maxBytes {
|
||||
t.Fatalf("%d bytes of report files kept, leaving no room "+
|
||||
"for the report", after)
|
||||
}
|
||||
|
||||
err = buf.Flush()
|
||||
@@ -258,8 +337,217 @@ func TestWrittenReportsKeepCounting(t *testing.T) {
|
||||
t.Fatalf("flush: %v", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
t.Fatal("the report files never filled the cap")
|
||||
// TestFailedDeletionStillCounts makes deleting the oldest report file
|
||||
// fail. It is still there, so it still takes room, and the next oldest
|
||||
// is deleted in its place.
|
||||
func TestFailedDeletionStillCounts(t *testing.T) {
|
||||
const fileBytes = 100
|
||||
|
||||
report := map[string]string{"id": "stuck"}
|
||||
dir := t.TempDir()
|
||||
stuck := reportFilePath(dir, 1)
|
||||
next := reportFilePath(dir, 2)
|
||||
newest := reportFilePath(dir, 3)
|
||||
|
||||
for _, path := range []string{stuck, next, newest} {
|
||||
writeBytes(t, path, fileBytes)
|
||||
}
|
||||
|
||||
t.Setenv("DATA_DIR", dir)
|
||||
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(3*fileBytes))
|
||||
|
||||
buf := startBuffer(t)
|
||||
|
||||
// A directory that is not empty cannot be deleted, even by root,
|
||||
// which the tests run as in the backend image.
|
||||
err := os.Remove(stuck)
|
||||
if err != nil {
|
||||
t.Fatalf("remove %s: %v", stuck, err)
|
||||
}
|
||||
|
||||
err = os.Mkdir(stuck, 0o750)
|
||||
if err != nil {
|
||||
t.Fatalf("make directory %s: %v", stuck, err)
|
||||
}
|
||||
|
||||
writeBytes(t, filepath.Join(stuck, "file"), 1)
|
||||
|
||||
err = buf.Append(report)
|
||||
if err != nil {
|
||||
t.Fatalf("report past the cap: %v", err)
|
||||
}
|
||||
|
||||
if exists(t, next) {
|
||||
t.Fatal("next oldest report file kept: the failed deletion " +
|
||||
"counted as making room")
|
||||
}
|
||||
|
||||
if !exists(t, newest) {
|
||||
t.Fatal("newest report file deleted, though deleting one made room")
|
||||
}
|
||||
}
|
||||
|
||||
// TestFileDeletedByHandFreesRoom deletes the oldest report file by
|
||||
// hand after start. When room is needed, its room counts as freed, so
|
||||
// no other file is deleted.
|
||||
func TestFileDeletedByHandFreesRoom(t *testing.T) {
|
||||
const fileBytes = 100
|
||||
|
||||
report := map[string]string{"id": "by-hand"}
|
||||
dir := t.TempDir()
|
||||
gone := reportFilePath(dir, 1)
|
||||
kept := reportFilePath(dir, 2)
|
||||
|
||||
writeBytes(t, gone, fileBytes)
|
||||
writeBytes(t, kept, fileBytes)
|
||||
|
||||
t.Setenv("DATA_DIR", dir)
|
||||
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(2*fileBytes))
|
||||
|
||||
buf := startBuffer(t)
|
||||
|
||||
err := os.Remove(gone)
|
||||
if err != nil {
|
||||
t.Fatalf("remove %s: %v", gone, err)
|
||||
}
|
||||
|
||||
err = buf.Append(report)
|
||||
if err != nil {
|
||||
t.Fatalf("report past the cap: %v", err)
|
||||
}
|
||||
|
||||
if !exists(t, kept) {
|
||||
t.Fatal("report file deleted, though the one deleted by hand " +
|
||||
"had made room")
|
||||
}
|
||||
}
|
||||
|
||||
// TestFileBeingWrittenIsNeverDeleted holds the write of one report file
|
||||
// open while a second write completes, then sends a report that needs
|
||||
// room. Deleting either file would make it, and the one being written is
|
||||
// the older, but only the complete one may be deleted. Once the first
|
||||
// write is complete, its file is deleted when room is needed.
|
||||
func TestFileBeingWrittenIsNeverDeleted(t *testing.T) {
|
||||
report := map[string]string{"id": "writing"}
|
||||
|
||||
t.Setenv("DATA_DIR", t.TempDir())
|
||||
// Room for two reports waiting to be written, but not for two
|
||||
// beside a report file.
|
||||
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(2*lineBytes(t, report)))
|
||||
|
||||
buf := startBuffer(t)
|
||||
|
||||
created := make(chan string)
|
||||
release := make(chan struct{})
|
||||
|
||||
buf.OnFileCreated(func(f *os.File) {
|
||||
created <- f.Name()
|
||||
|
||||
<-release
|
||||
})
|
||||
|
||||
err := buf.Append(report)
|
||||
if err != nil {
|
||||
t.Fatalf("first report: %v", err)
|
||||
}
|
||||
|
||||
flushed := make(chan error)
|
||||
|
||||
go func() { flushed <- buf.Flush() }()
|
||||
|
||||
writing := <-created
|
||||
|
||||
// Only the first write is held; the second goes through, and
|
||||
// so do the writes after it, the final one at stop included.
|
||||
var complete string
|
||||
|
||||
buf.OnFileCreated(func(f *os.File) { complete = f.Name() })
|
||||
|
||||
// Errorf, not Fatalf, until the first write is released, so that a
|
||||
// failure here does not leave it held.
|
||||
err = buf.Append(report)
|
||||
if err != nil {
|
||||
t.Errorf("second report: %v", err)
|
||||
}
|
||||
|
||||
err = buf.Flush()
|
||||
if err != nil {
|
||||
t.Errorf("flush of the second report: %v", err)
|
||||
}
|
||||
|
||||
err = buf.Append(report)
|
||||
if err != nil {
|
||||
t.Errorf("report that needs room: %v", err)
|
||||
}
|
||||
|
||||
if !exists(t, writing) {
|
||||
t.Error("report file deleted while it was being written")
|
||||
}
|
||||
|
||||
if exists(t, complete) {
|
||||
t.Error("complete report file kept, though room was needed")
|
||||
}
|
||||
|
||||
close(release)
|
||||
|
||||
err = <-flushed
|
||||
if err != nil {
|
||||
t.Fatalf("flush of the first report: %v", err)
|
||||
}
|
||||
|
||||
err = buf.Append(report)
|
||||
if err != nil {
|
||||
t.Fatalf("report after the first write was complete: %v", err)
|
||||
}
|
||||
|
||||
if exists(t, writing) {
|
||||
t.Fatal("complete report file kept when room was needed")
|
||||
}
|
||||
}
|
||||
|
||||
// TestFailedWriteStopsCounting makes a write fail once its file is
|
||||
// created. Its reports are lost, so they stop counting, and the part of
|
||||
// the file written is removed, so it takes no room.
|
||||
func TestFailedWriteStopsCounting(t *testing.T) {
|
||||
report := map[string]string{"id": "failed"}
|
||||
|
||||
t.Setenv("DATA_DIR", t.TempDir())
|
||||
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(lineBytes(t, report)))
|
||||
|
||||
buf := startBuffer(t)
|
||||
|
||||
var failed string
|
||||
|
||||
// Closing the file under the write makes the write fail.
|
||||
buf.OnFileCreated(func(f *os.File) {
|
||||
failed = f.Name()
|
||||
|
||||
_ = f.Close()
|
||||
})
|
||||
|
||||
err := buf.Append(report)
|
||||
if err != nil {
|
||||
t.Fatalf("report that fills the cap exactly: %v", err)
|
||||
}
|
||||
|
||||
err = buf.Flush()
|
||||
if err == nil {
|
||||
t.Fatal("flush succeeded, though its file was closed under it")
|
||||
}
|
||||
|
||||
if exists(t, failed) {
|
||||
t.Fatal("file of the failed write kept")
|
||||
}
|
||||
|
||||
// Writes from here on, the final one at stop included, succeed.
|
||||
buf.OnFileCreated(func(*os.File) {})
|
||||
|
||||
err = buf.Append(report)
|
||||
if err != nil {
|
||||
t.Fatalf("report after the failed write: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestConcurrentAppendsStopAtCap appends from many goroutines at once
|
||||
@@ -334,7 +622,11 @@ func TestTwoFlushesInOneMillisecond(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
files := readReportFiles(t, dir)
|
||||
files, err := readReportFiles(dir)
|
||||
if err != nil {
|
||||
t.Fatalf("read report files: %v", err)
|
||||
}
|
||||
|
||||
if len(files) != flushes {
|
||||
t.Fatalf("%d report files after %d flushes", len(files), flushes)
|
||||
}
|
||||
@@ -347,6 +639,88 @@ func TestTwoFlushesInOneMillisecond(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestReportFileHoldsTheLinesAppended flushes three reports and reads
|
||||
// their file back: it must decompress to exactly their JSON lines, in
|
||||
// the order they were appended.
|
||||
func TestReportFileHoldsTheLinesAppended(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
t.Setenv("DATA_DIR", dir)
|
||||
|
||||
buf := startBuffer(t)
|
||||
|
||||
for _, id := range []int{1, 2, 3} {
|
||||
err := buf.Append(map[string]int{"id": id})
|
||||
if err != nil {
|
||||
t.Fatalf("append report %d: %v", id, err)
|
||||
}
|
||||
}
|
||||
|
||||
err := buf.Flush()
|
||||
if err != nil {
|
||||
t.Fatalf("flush: %v", err)
|
||||
}
|
||||
|
||||
files, err := readReportFiles(dir)
|
||||
if err != nil {
|
||||
t.Fatalf("read report files: %v", err)
|
||||
}
|
||||
|
||||
want := `{"id":1}` + "\n" + `{"id":2}` + "\n" + `{"id":3}` + "\n"
|
||||
if len(files) != 1 || files[0] != want {
|
||||
t.Fatalf("report files = %q, want one holding %q", files, want)
|
||||
}
|
||||
}
|
||||
|
||||
// TestFlushAtSizeThreshold appends reports until the buffer holds
|
||||
// FlushSizeThreshold bytes. The append that gets it there must write
|
||||
// them all to one report file, with no call to Flush and the periodic
|
||||
// flush a minute away, and no earlier append may write one.
|
||||
func TestFlushAtSizeThreshold(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
t.Setenv("DATA_DIR", dir)
|
||||
|
||||
buf := startBuffer(t)
|
||||
|
||||
// Large reports, so the threshold takes a few hundred appends.
|
||||
pad := strings.Repeat("a", 64<<10)
|
||||
|
||||
var appended strings.Builder
|
||||
|
||||
for id := 0; appended.Len() < reportbuf.FlushSizeThreshold; id++ {
|
||||
report := map[string]any{"id": id, "pad": pad}
|
||||
|
||||
err := buf.Append(report)
|
||||
if err != nil {
|
||||
t.Fatalf("append report %d: %v", id, err)
|
||||
}
|
||||
|
||||
line, err := json.Marshal(report)
|
||||
if err != nil {
|
||||
t.Fatalf("marshal report %d: %v", id, err)
|
||||
}
|
||||
|
||||
appended.Write(line)
|
||||
appended.WriteByte('\n')
|
||||
}
|
||||
|
||||
// Append writes the file in the background, so wait for it.
|
||||
deadline := time.Now().Add(10 * time.Second)
|
||||
|
||||
for {
|
||||
files, err := readReportFiles(dir)
|
||||
if err == nil && len(files) == 1 && files[0] == appended.String() {
|
||||
return
|
||||
}
|
||||
|
||||
if time.Now().After(deadline) {
|
||||
t.Fatalf("%d report files (error: %v), want one holding the "+
|
||||
"%d bytes appended", len(files), err, appended.Len())
|
||||
}
|
||||
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
}
|
||||
}
|
||||
|
||||
// reportFilesBytes returns the total size of the report files in dir.
|
||||
func reportFilesBytes(t *testing.T, dir string) int64 {
|
||||
t.Helper()
|
||||
@@ -371,20 +745,19 @@ func reportFilesBytes(t *testing.T, dir string) int64 {
|
||||
}
|
||||
|
||||
// readReportFiles returns the decompressed contents of each report
|
||||
// file in dir.
|
||||
func readReportFiles(t *testing.T, dir string) []string {
|
||||
t.Helper()
|
||||
|
||||
// file in dir. A file still being written does not decompress, so it
|
||||
// gives an error.
|
||||
func readReportFiles(dir string) ([]string, error) {
|
||||
files := os.DirFS(dir)
|
||||
|
||||
names, err := fs.Glob(files, "reports-*.jsonl.zst")
|
||||
if err != nil {
|
||||
t.Fatalf("list report files: %v", err)
|
||||
return nil, fmt.Errorf("list report files: %w", err)
|
||||
}
|
||||
|
||||
dec, err := zstd.NewReader(nil)
|
||||
if err != nil {
|
||||
t.Fatalf("create zstd decoder: %v", err)
|
||||
return nil, fmt.Errorf("create zstd decoder: %w", err)
|
||||
}
|
||||
defer dec.Close()
|
||||
|
||||
@@ -393,18 +766,26 @@ func readReportFiles(t *testing.T, dir string) []string {
|
||||
for _, name := range names {
|
||||
compressed, readErr := fs.ReadFile(files, name)
|
||||
if readErr != nil {
|
||||
t.Fatalf("read %s: %v", name, readErr)
|
||||
return nil, fmt.Errorf("read %s: %w", name, readErr)
|
||||
}
|
||||
|
||||
data, decErr := dec.DecodeAll(compressed, nil)
|
||||
if decErr != nil {
|
||||
t.Fatalf("decompress %s: %v", name, decErr)
|
||||
return nil, fmt.Errorf("decompress %s: %w", name, decErr)
|
||||
}
|
||||
|
||||
contents = append(contents, string(data))
|
||||
}
|
||||
|
||||
return contents
|
||||
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) {
|
||||
@@ -416,6 +797,21 @@ func writeBytes(t *testing.T, path string, n int) {
|
||||
}
|
||||
}
|
||||
|
||||
func exists(t *testing.T, path string) bool {
|
||||
t.Helper()
|
||||
|
||||
_, err := os.Stat(path)
|
||||
if errors.Is(err, fs.ErrNotExist) {
|
||||
return false
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
t.Fatalf("stat %s: %v", path, err)
|
||||
}
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
func hasReportFile(t *testing.T, dir string) bool {
|
||||
t.Helper()
|
||||
|
||||
@@ -439,3 +835,85 @@ func hasReportFile(t *testing.T, dir string) bool {
|
||||
|
||||
return false
|
||||
}
|
||||
|
||||
// TestFilesDeletedOldestFirstWhenWritesOverlap holds the write of an
|
||||
// older report file open until a newer one's write completes, then
|
||||
// releases it. When room is needed, the older file is deleted first,
|
||||
// though its write was the last to complete.
|
||||
func TestFilesDeletedOldestFirstWhenWritesOverlap(t *testing.T) {
|
||||
const maxBytes = 1000
|
||||
|
||||
report := map[string]string{"id": "overlap"}
|
||||
|
||||
t.Setenv("DATA_DIR", t.TempDir())
|
||||
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(maxBytes))
|
||||
|
||||
buf := startBuffer(t)
|
||||
|
||||
created := make(chan string)
|
||||
release := make(chan struct{})
|
||||
|
||||
buf.OnFileCreated(func(f *os.File) {
|
||||
created <- f.Name()
|
||||
|
||||
<-release
|
||||
})
|
||||
|
||||
err := buf.Append(report)
|
||||
if err != nil {
|
||||
t.Fatalf("older report: %v", err)
|
||||
}
|
||||
|
||||
flushed := make(chan error)
|
||||
|
||||
go func() { flushed <- buf.Flush() }()
|
||||
|
||||
older := <-created
|
||||
|
||||
// Only the older write is held; the newer one goes through, and
|
||||
// so do the writes after it, the final one at stop included.
|
||||
var newer string
|
||||
|
||||
buf.OnFileCreated(func(f *os.File) { newer = f.Name() })
|
||||
|
||||
// Errorf, not Fatalf, until the older write is released, so that a
|
||||
// failure here does not leave it held.
|
||||
err = buf.Append(report)
|
||||
if err != nil {
|
||||
t.Errorf("newer report: %v", err)
|
||||
}
|
||||
|
||||
err = buf.Flush()
|
||||
if err != nil {
|
||||
t.Errorf("flush of the newer report: %v", err)
|
||||
}
|
||||
|
||||
close(release)
|
||||
|
||||
err = <-flushed
|
||||
if err != nil {
|
||||
t.Fatalf("flush of the older report: %v", err)
|
||||
}
|
||||
|
||||
info, err := os.Stat(newer)
|
||||
if err != nil {
|
||||
t.Fatalf("stat %s: %v", newer, err)
|
||||
}
|
||||
|
||||
// A report that fits beside the newer file alone, so deleting the
|
||||
// older one makes exactly the room it needs.
|
||||
pad := maxBytes - int(info.Size()) - lineBytes(t, map[string]string{"id": ""})
|
||||
|
||||
err = buf.Append(map[string]string{"id": strings.Repeat("a", pad)})
|
||||
if err != nil {
|
||||
t.Fatalf("report that needs room: %v", err)
|
||||
}
|
||||
|
||||
if exists(t, older) {
|
||||
t.Fatal("older report file kept when room was needed")
|
||||
}
|
||||
|
||||
if !exists(t, newer) {
|
||||
t.Fatal("newer report file deleted before the older one")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"errors"
|
||||
"net"
|
||||
"net/http"
|
||||
"runtime"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
@@ -53,7 +54,7 @@ func (s *Server) listenAndServe() {
|
||||
s.log.Info("http begin listen",
|
||||
"listenaddr", s.httpServer.Addr,
|
||||
"version", s.params.Globals.Version,
|
||||
"buildarch", s.params.Globals.Buildarch,
|
||||
"arch", runtime.GOARCH,
|
||||
)
|
||||
|
||||
err := s.httpServer.ListenAndServe()
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
#!/bin/sh
|
||||
# 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
|
||||
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
||||
@@ -14,7 +14,7 @@ main() {
|
||||
version="${VERSION:-$(git describe --always --dirty 2>/dev/null || echo dev)}"
|
||||
|
||||
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/
|
||||
}
|
||||
|
||||
|
||||
+20
-2
@@ -14,16 +14,34 @@ set -eu
|
||||
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
||||
|
||||
# The sha256 of the org standard .golangci.yml. When that file changes in
|
||||
# sneak/prompts and is copied here again, this changes with it.
|
||||
GOLANGCI_CONFIG_SHA256="a79b63a254602a5318db5d0e9a06bc71b84bf0c1d896305229d8bfed1d1b1776"
|
||||
|
||||
main() {
|
||||
cd "$ROOT"
|
||||
if [ ! -f .golangci.yml ]; then
|
||||
echo "backend/.golangci.yml is missing. Copy the org standard verbatim" >&2
|
||||
echo "from https://git.eeqj.de/sneak/prompts/raw/branch/main/.golangci.yml" >&2
|
||||
exit 1
|
||||
fi
|
||||
actual="$(sha256sum .golangci.yml | cut -d' ' -f1)"
|
||||
if [ -z "$actual" ]; then
|
||||
echo "sha256sum is missing or printed no hash, so" >&2
|
||||
echo "backend/.golangci.yml could not be checked." >&2
|
||||
exit 1
|
||||
fi
|
||||
if [ "$actual" != "$GOLANGCI_CONFIG_SHA256" ]; then
|
||||
echo ".golangci.yml has drifted from the org standard." >&2
|
||||
echo "backend/.golangci.yml does not match GOLANGCI_CONFIG_SHA256" >&2
|
||||
echo "in backend/script/lint." >&2
|
||||
echo " expected $GOLANGCI_CONFIG_SHA256" >&2
|
||||
echo " actual $actual" >&2
|
||||
echo "Restore it verbatim from sneak/prompts; do not edit it." >&2
|
||||
echo "Compare it with the org standard," >&2
|
||||
echo "https://git.eeqj.de/sneak/prompts/raw/branch/main/.golangci.yml" >&2
|
||||
echo "- If they differ, it was edited here: restore the org standard" >&2
|
||||
echo " verbatim. Do not edit it." >&2
|
||||
echo "- If they are the same, the org standard changed: set" >&2
|
||||
echo " GOLANGCI_CONFIG_SHA256 in backend/script/lint to the actual hash." >&2
|
||||
exit 1
|
||||
fi
|
||||
golangci-lint run ./...
|
||||
|
||||
+10
-2
@@ -1,12 +1,20 @@
|
||||
#!/bin/sh
|
||||
# script/test: run the backend test suite.
|
||||
# script/test: run the backend test suite with the race detector and
|
||||
# coverage. Go's own -timeout bounds the tests and not their compile,
|
||||
# so a cold build cache cannot fail it. The race detector needs cgo,
|
||||
# and so a C compiler. If the tests fail, they run again with -v for
|
||||
# the details, and the script fails even if that run passes.
|
||||
set -eu
|
||||
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
||||
|
||||
main() {
|
||||
cd "$ROOT"
|
||||
timeout 30 go test ./...
|
||||
go test -timeout 30s -race -cover ./... || {
|
||||
echo "--- Rerunning with -v for details ---"
|
||||
go test -timeout 30s -race -v ./...
|
||||
exit 1
|
||||
}
|
||||
}
|
||||
|
||||
main "$@"
|
||||
|
||||
+6
-19
@@ -63,26 +63,13 @@ done > /etc/nginx/trusted-proxies.conf
|
||||
|
||||
# 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
|
||||
# uid. Both are given to the netwatch user here, with the mode the
|
||||
# server gives a directory it creates, so the host directory needs no
|
||||
# preparing.
|
||||
#
|
||||
# chown and chmod, run as root, change whatever a symbolic link on the
|
||||
# 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.
|
||||
# uid. Here, as root, netwatch-server prepare-data-dir creates DATA_DIR
|
||||
# and gives /data and everything in it to the netwatch user, so the
|
||||
# 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
|
||||
# nothing outside /data, whatever symbolic links it meets there.
|
||||
export DATA_DIR="${DATA_DIR:-/data/reports}"
|
||||
mkdir -p "$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
|
||||
netwatch-server prepare-data-dir "$DATA_DIR" || exit 1
|
||||
|
||||
# A stop signal is only noted here; the loop below acts on it.
|
||||
stop_requested=""
|
||||
|
||||
@@ -9,6 +9,9 @@
|
||||
type="image/svg+xml"
|
||||
href="data:image/svg+xml,<svg xmlns='http://www.w3.org/2000/svg' viewBox='0 0 100 100'><text y='.9em' font-size='90'>📡</text></svg>"
|
||||
/>
|
||||
<!-- Linked here, not imported by src/main.js, so the unit tests can
|
||||
import that module in Node, which cannot import CSS. -->
|
||||
<link rel="stylesheet" href="/src/styles.css" />
|
||||
</head>
|
||||
<body class="bg-gray-900 text-white min-h-screen">
|
||||
<div id="app"></div>
|
||||
|
||||
+10
-3
@@ -64,7 +64,8 @@ detect_pkgmgr() {
|
||||
fi
|
||||
}
|
||||
|
||||
# pkg_install <nix-attr> <apt-pkg> <brew-formula> <apk-pkg>
|
||||
# pkg_install <nix-attr> <apt-pkgs> <brew-formula> <apk-pkgs>: the apt
|
||||
# and apk arguments may each list several packages, separated by spaces.
|
||||
pkg_install() {
|
||||
detect_pkgmgr
|
||||
case "$PKGMGR" in
|
||||
@@ -74,10 +75,10 @@ pkg_install() {
|
||||
$SUDO env DEBIAN_FRONTEND=noninteractive apt-get update
|
||||
APT_UPDATED=1
|
||||
fi
|
||||
$SUDO env DEBIAN_FRONTEND=noninteractive apt-get install -y "$2"
|
||||
$SUDO env DEBIAN_FRONTEND=noninteractive apt-get install -y $2
|
||||
;;
|
||||
brew) brew install "$3" ;;
|
||||
apk) apk add --no-cache "$4" ;;
|
||||
apk) apk add --no-cache $4 ;;
|
||||
esac
|
||||
}
|
||||
|
||||
@@ -253,6 +254,12 @@ main() {
|
||||
|
||||
if missing make; then pkg_install gnumake make make make; fi
|
||||
if missing git; then pkg_install git git git git; fi
|
||||
# The race detector in make test needs cgo, which Go turns on only
|
||||
# when it finds its C compiler, gcc on Linux. apt and apk ship the C
|
||||
# library headers apart from gcc.
|
||||
if missing gcc; then
|
||||
pkg_install gcc "gcc libc6-dev" gcc "gcc musl-dev"
|
||||
fi
|
||||
|
||||
ensure_node
|
||||
ensure_yarn
|
||||
|
||||
+2
-3
@@ -16,9 +16,8 @@ main() {
|
||||
"$SCRIPT_DIR/check"
|
||||
# Own line: a failing command substitution inside an argument does
|
||||
# not trip `set -e`, so the inline form degrades silently to an
|
||||
# empty constant. VERSION is computed here because .dockerignore
|
||||
# excludes .git, so `git describe` in a build stage yields an empty
|
||||
# version without failing.
|
||||
# empty constant. The VERSION build argument takes precedence over
|
||||
# the version a build stage derives from the .git in the context.
|
||||
version="$(git describe --tags --always --dirty 2>/dev/null || true)"
|
||||
[ -n "$version" ] || version="unknown"
|
||||
docker build --no-cache \
|
||||
|
||||
+2
-3
@@ -12,9 +12,8 @@ main() {
|
||||
cd "$ROOT"
|
||||
# Own line: a failing command substitution inside an argument does
|
||||
# not trip `set -e`, so the inline form degrades silently to an
|
||||
# empty constant. VERSION is computed here because .dockerignore
|
||||
# excludes .git, so `git describe` in a build stage yields an empty
|
||||
# version without failing.
|
||||
# empty constant. The VERSION build argument takes precedence over
|
||||
# the version a build stage derives from the .git in the context.
|
||||
version="$(git describe --tags --always --dirty 2>/dev/null || true)"
|
||||
[ -n "$version" ] || version="unknown"
|
||||
docker build --no-cache \
|
||||
|
||||
@@ -1,13 +1,14 @@
|
||||
#!/bin/sh
|
||||
# script/frontend-test: run the frontend test suite. The frontend has no
|
||||
# unit tests; the production build serves as the test (fails on broken
|
||||
# code).
|
||||
# script/frontend-test: run the frontend test suite: the unit tests in
|
||||
# test/unit/ with Node's built-in test runner, then the production
|
||||
# build, which fails on broken code.
|
||||
set -eu
|
||||
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
||||
|
||||
main() {
|
||||
cd "$ROOT"
|
||||
timeout 30 node --test test/unit/*.test.js
|
||||
timeout 30 yarn build
|
||||
}
|
||||
|
||||
|
||||
+6
-3
@@ -1,14 +1,17 @@
|
||||
#!/bin/sh
|
||||
# script/test: run the test suite for the whole repo: the frontend at
|
||||
# the repo root, then the Go backend in backend/. Both halves together
|
||||
# get 30 seconds; each also keeps its own limit for the Dockerfiles.
|
||||
# the repo root, then the Go backend in backend/. Each half has its own
|
||||
# 30-second limit, and there is none around both: from a cold Go build
|
||||
# cache, compiling the backend's tests with the race detector can take
|
||||
# 30 seconds on its own, and Go's -timeout leaves the compile out.
|
||||
set -eu
|
||||
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
||||
|
||||
main() {
|
||||
cd "$ROOT"
|
||||
timeout 30 sh -c 'script/frontend-test && backend/script/test'
|
||||
script/frontend-test
|
||||
backend/script/test
|
||||
}
|
||||
|
||||
main "$@"
|
||||
|
||||
+54
-29
@@ -1,14 +1,14 @@
|
||||
import "./styles.css";
|
||||
|
||||
// --- Configuration -----------------------------------------------------------
|
||||
|
||||
// Timing, axis labels, and display constants. Latency above maxLatency is
|
||||
// clamped to "unreachable". The sparkline Y-axis is capped at
|
||||
// Timing, axis labels, and display constants. A target check times out
|
||||
// after requestTimeout, 80% of updateInterval, so a round's checks have
|
||||
// all finished before the next round is due; latency above maxLatency is
|
||||
// recorded as a timeout. The sparkline Y-axis is capped at
|
||||
// graphMaxLatency — values above it pin to the top of the chart but still
|
||||
// display their real value in the latency figure. The history buffer holds
|
||||
// maxHistoryPoints samples (historyDuration / updateInterval).
|
||||
// reportInterval is how often collected samples are POSTed to the backend.
|
||||
const CONFIG = {
|
||||
export const CONFIG = {
|
||||
updateInterval: 3000,
|
||||
maxHistoryPoints: 100,
|
||||
reportInterval: 60000,
|
||||
@@ -16,7 +16,7 @@ const CONFIG = {
|
||||
return (this.maxHistoryPoints * this.updateInterval) / 1000;
|
||||
},
|
||||
get requestTimeout() {
|
||||
return Math.min(this.updateInterval - 100, 3000);
|
||||
return this.updateInterval * 0.8;
|
||||
},
|
||||
get maxLatency() {
|
||||
return this.requestTimeout;
|
||||
@@ -503,12 +503,16 @@ class Reporter {
|
||||
|
||||
// --- Latency Measurement -----------------------------------------------------
|
||||
|
||||
async function measureLatency(url) {
|
||||
// Checks one target. The check times out after CONFIG.requestTimeout; the
|
||||
// caller can give it up sooner through the optional signal, which also ends
|
||||
// it as a timeout.
|
||||
export async function measureLatency(url, signal) {
|
||||
const controller = new AbortController();
|
||||
const timeoutId = setTimeout(
|
||||
() => controller.abort(),
|
||||
CONFIG.requestTimeout,
|
||||
);
|
||||
signal?.addEventListener("abort", () => controller.abort());
|
||||
|
||||
const targetUrl = new URL(url);
|
||||
targetUrl.searchParams.set("_cb", Date.now().toString());
|
||||
@@ -1101,7 +1105,7 @@ function sortAndRebuildWAN(state) {
|
||||
|
||||
// --- Main Loop ---------------------------------------------------------------
|
||||
|
||||
async function tick(state, onOffline) {
|
||||
async function tick(state, signal, onOffline) {
|
||||
const ts = Date.now();
|
||||
|
||||
if (state.paused) {
|
||||
@@ -1123,11 +1127,12 @@ async function tick(state, onOffline) {
|
||||
log.debug(`Tick #${state.tickCount + 1} started`);
|
||||
|
||||
const results = await Promise.all(
|
||||
state.allHosts.map((h) => measureLatency(h.url)),
|
||||
state.allHosts.map((h) => measureLatency(h.url, signal)),
|
||||
);
|
||||
|
||||
// User may have paused while awaiting results — discard them
|
||||
if (state.paused) return;
|
||||
// User may have paused, or the next round may have given up this
|
||||
// one's checks, while awaiting results — discard them
|
||||
if (state.paused || signal.aborted) return;
|
||||
|
||||
state.tickCount++;
|
||||
|
||||
@@ -1168,9 +1173,10 @@ async function tick(state, onOffline) {
|
||||
|
||||
// --- Recovery Probe ----------------------------------------------------------
|
||||
|
||||
// When offline, rapidly poll 4 random WAN hosts every 500ms. As soon as any
|
||||
// responds, stop probing and fire a normal tick to refresh all hosts.
|
||||
function startRecoveryProbe(state, triggerTick) {
|
||||
// When offline, check 4 random WAN hosts every 500ms, giving up the checks
|
||||
// started 500ms before, so at most 4 are ever waiting. As soon as one
|
||||
// answers, stop probing and start a new round at once.
|
||||
function startRecoveryProbe(state, startRounds) {
|
||||
if (state._recoveryProbeId) return; // already running
|
||||
const candidates = [...state.wan];
|
||||
for (let i = candidates.length - 1; i > 0; i--) {
|
||||
@@ -1181,15 +1187,18 @@ function startRecoveryProbe(state, triggerTick) {
|
||||
log.notice(
|
||||
`Recovery probe started (${canaries.map((h) => h.name).join(", ")})`,
|
||||
);
|
||||
state._recoveryProbeId = setInterval(async () => {
|
||||
state._recoveryProbeId = setInterval(() => {
|
||||
if (state.paused) return;
|
||||
const results = await Promise.all(
|
||||
canaries.map((h) => measureLatency(h.url)),
|
||||
);
|
||||
if (results.some((r) => r.error === null)) {
|
||||
log.notice("Recovery probe: connectivity detected");
|
||||
stopRecoveryProbe(state);
|
||||
triggerTick();
|
||||
state._recoveryProbeChecks?.abort();
|
||||
const checks = new AbortController();
|
||||
state._recoveryProbeChecks = checks;
|
||||
for (const host of canaries) {
|
||||
measureLatency(host.url, checks.signal).then((r) => {
|
||||
if (r.error !== null || checks.signal.aborted) return;
|
||||
log.notice("Recovery probe: connectivity detected");
|
||||
stopRecoveryProbe(state);
|
||||
startRounds();
|
||||
});
|
||||
}
|
||||
}, 500);
|
||||
}
|
||||
@@ -1198,6 +1207,7 @@ function stopRecoveryProbe(state) {
|
||||
if (state._recoveryProbeId) {
|
||||
clearInterval(state._recoveryProbeId);
|
||||
state._recoveryProbeId = null;
|
||||
state._recoveryProbeChecks?.abort();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1374,18 +1384,34 @@ async function init() {
|
||||
updateClocks();
|
||||
setInterval(updateClocks, 1000);
|
||||
|
||||
// Rounds never overlap: a round first gives up the last round's checks
|
||||
// if they are still waiting, and the last round then records nothing.
|
||||
// At a steady interval they never are, as they time out at 80% of it;
|
||||
// they can be when a round starts early, after an interval change or
|
||||
// when the recovery probe finds a target answering.
|
||||
let roundChecks = new AbortController();
|
||||
function doTick() {
|
||||
tick(state, () => startRecoveryProbe(state, doTick));
|
||||
roundChecks.abort();
|
||||
roundChecks = new AbortController();
|
||||
tick(state, roundChecks.signal, () =>
|
||||
startRecoveryProbe(state, startRounds),
|
||||
);
|
||||
}
|
||||
|
||||
doTick();
|
||||
let tickIntervalId = setInterval(doTick, CONFIG.updateInterval);
|
||||
// Starts a round now and then one every CONFIG.updateInterval.
|
||||
let tickIntervalId;
|
||||
function startRounds() {
|
||||
clearInterval(tickIntervalId);
|
||||
doTick();
|
||||
tickIntervalId = setInterval(doTick, CONFIG.updateInterval);
|
||||
}
|
||||
|
||||
startRounds();
|
||||
|
||||
document
|
||||
.getElementById("interval-select")
|
||||
.addEventListener("change", (e) => {
|
||||
const newInterval = parseInt(e.target.value, 10);
|
||||
clearInterval(tickIntervalId);
|
||||
CONFIG.updateInterval = newInterval;
|
||||
log.notice(
|
||||
`Interval changed to ${humanDuration(newInterval / 1000)}, history reset`,
|
||||
@@ -1434,8 +1460,7 @@ async function init() {
|
||||
|
||||
// Start immediately with new interval
|
||||
stopRecoveryProbe(state);
|
||||
doTick();
|
||||
tickIntervalId = setInterval(doTick, CONFIG.updateInterval);
|
||||
startRounds();
|
||||
});
|
||||
|
||||
window.addEventListener("resize", () => handleResize(state));
|
||||
@@ -1444,7 +1469,7 @@ async function init() {
|
||||
|
||||
// Bootstrap only when loaded as the page: a real DOM containing the #app
|
||||
// mount point this module renders into. Importing the module in a unit test
|
||||
// (which has no #app) runs nothing, so buildReport can be tested in isolation.
|
||||
// (which has no #app) runs nothing, so its exports can be tested in isolation.
|
||||
if (typeof document !== "undefined" && document.getElementById("app")) {
|
||||
if (document.readyState === "loading") {
|
||||
document.addEventListener("DOMContentLoaded", init);
|
||||
|
||||
@@ -0,0 +1,66 @@
|
||||
// Unit tests for src/main.js, run by script/frontend-test with Node's
|
||||
// built-in test runner. Importing the module does not start the page.
|
||||
|
||||
import { test } from "node:test";
|
||||
import assert from "node:assert/strict";
|
||||
import { CONFIG, measureLatency } from "../../src/main.js";
|
||||
|
||||
// measureLatency writes timeouts to the debug log, which looks for its
|
||||
// panel in the page. There is no page here.
|
||||
globalThis.document = { getElementById: () => null };
|
||||
|
||||
// Mocks the clock for test t, so that a check lasting seconds takes no real
|
||||
// time, and replaces fetch with a target that answers after answerAfter
|
||||
// milliseconds of that clock, or never when answerAfter is Infinity. Both
|
||||
// are restored when the test ends.
|
||||
function mockTarget(t, answerAfter) {
|
||||
t.mock.timers.enable({ apis: ["setTimeout", "Date"] });
|
||||
t.mock.method(performance, "now", () => Date.now());
|
||||
t.mock.method(
|
||||
globalThis,
|
||||
"fetch",
|
||||
(url, { signal }) =>
|
||||
new Promise((resolve, reject) => {
|
||||
if (answerAfter !== Infinity) setTimeout(resolve, answerAfter);
|
||||
signal.addEventListener("abort", () => reject(signal.reason));
|
||||
}),
|
||||
);
|
||||
}
|
||||
|
||||
// The result of check if it has ended, otherwise "still waiting".
|
||||
function settled(check) {
|
||||
return Promise.race([
|
||||
check,
|
||||
new Promise((resolve) => setImmediate(resolve, "still waiting")),
|
||||
]);
|
||||
}
|
||||
|
||||
for (const interval of [10000, 30000]) {
|
||||
const timeout = interval * 0.8;
|
||||
// Over 3 seconds, which the timeout used to be capped at.
|
||||
const slowAnswer = timeout - 1000;
|
||||
|
||||
test(`at a ${interval}ms interval, an answer after ${slowAnswer}ms is recorded with its real time`, async (t) => {
|
||||
CONFIG.updateInterval = interval;
|
||||
mockTarget(t, slowAnswer);
|
||||
const check = measureLatency("https://target.test");
|
||||
t.mock.timers.tick(slowAnswer);
|
||||
assert.deepEqual(await settled(check), {
|
||||
latency: slowAnswer,
|
||||
error: null,
|
||||
});
|
||||
});
|
||||
|
||||
test(`at a ${interval}ms interval, a target that never answers is recorded as a timeout after ${timeout}ms`, async (t) => {
|
||||
CONFIG.updateInterval = interval;
|
||||
mockTarget(t, Infinity);
|
||||
const check = measureLatency("https://target.test");
|
||||
t.mock.timers.tick(timeout - 1);
|
||||
assert.equal(await settled(check), "still waiting");
|
||||
t.mock.timers.tick(1);
|
||||
assert.deepEqual(await settled(check), {
|
||||
latency: null,
|
||||
error: "timeout",
|
||||
});
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user