4 Commits
Author SHA1 Message Date
sneak a2dd439812 Delete the oldest report files to stay under the size cap (closes #54)
check / check (push) Successful in 1m28s
When a report would take the report files past DATA_DIR_MAX_BYTES,
reportbuf now deletes the oldest report files until it fits, and does
the same at start when files left by an earlier run are already past
it. A file joins the files that may be deleted only once it is
completely written, so a file still being written is never deleted. A
report is refused with 507 only when the reports waiting to be written
fill the cap on their own, and then no file is deleted. The reports of
a failed write stop counting, and the part of its file written is
removed. A file whose deletion fails keeps counting; one already
deleted by hand counts as freed.

Model: opus-5-5
2026-10-03 14:12:07 +00:00
clawbot 0ab6418e3d Target check timeout is 80% of the refresh interval (closes #78)
check / check (push) Successful in 1m55s
Each target check now times out after 80% of the refresh interval, 24
seconds at 30 seconds, where it was capped at 3 seconds, so slow, far
targets are recorded with their real time. Rounds never overlap: a
round gives up the last round's checks if they are still waiting, which
happens only when a round starts early, after an interval change or
when the recovery probe finds a target answering. The probe still
checks every half second, giving up its previous checks, so they do not
pile up.

The frontend has its first unit tests, run by script/frontend-test with
Node's built-in test runner on a mocked clock. index.html now links
src/styles.css, which src/main.js imported, since Node cannot import
CSS.

Model: opus-5-5
2026-10-03 15:36:47 +02:00
clawbot 4ce0814b14 Stamp the git tag or short commit in a plain docker build (closes #86)
check / check (push) Successful in 1m36s
The builder stage now has git and takes the version from the VERSION
build argument when one is given, otherwise from `git describe --tags
--always` of the repo's .git, copied to /git so go build does not see
it. A version that still comes out empty, dev or unknown fails the
build. ARG VERSION has no default. .dockerignore keeps .git/config out
of the build context, and the comments in script/docker and
script/cibuild no longer say .git is excluded.

Model: opus-5-5
2026-10-02 04:26:40 +02:00
clawbot e6d6815ecb Remove Buildarch: read the architecture at run time (closes #83)
check / check (push) Successful in 1m37s
The architecture is no longer passed in at build time. The Buildarch
variable and field are gone from main and globals, script/build no
longer stamps it in with -X, and the startup and listen log lines
report runtime.GOARCH under the key "arch". The Dockerfile comment and
backend/README.md no longer describe an architecture being stamped in.

Model: opus-5-5
2026-10-02 01:05:09 +02:00
21 changed files with 708 additions and 170 deletions
+3
View File
@@ -4,3 +4,6 @@ tmp
.DS_Store .DS_Store
*.log *.log
.claude .claude
# .git is sent so the build can stamp the version, without its config.
.git/config
+20 -6
View File
@@ -20,7 +20,7 @@ RUN make lint
# golang:1.25-alpine (2026-02-27) # golang:1.25-alpine (2026-02-27)
FROM golang:1.25-alpine@sha256:f6751d823c26342f9506c03797d2527668d095b0a15f1862cddb4d927a7a4ced AS builder FROM golang:1.25-alpine@sha256:f6751d823c26342f9506c03797d2527668d095b0a15f1862cddb4d927a7a4ced AS builder
RUN apk add --no-cache make RUN apk add --no-cache git make
WORKDIR /src WORKDIR /src
@@ -37,11 +37,24 @@ RUN make test
# make build is a shim around backend/script/build, the one definition # make build is a shim around backend/script/build, the one definition
# of the build command: # of the build command:
# CGO_ENABLED=0 go build -trimpath -ldflags "-s -w -X main.Version=... -X main.Buildarch=..." # CGO_ENABLED=0 go build -trimpath -ldflags "-s -w -X main.Version=..."
# That script reads VERSION from the environment, so it is handed over # That script reads VERSION from the environment, so it is handed over
# there rather than as a make variable. # there rather than as a make variable.
ARG VERSION=dev #
RUN VERSION="${VERSION}" make build # The version is the VERSION build argument when one is given, otherwise
# `git describe --tags --always` of the repo's .git: the tag on a tagged
# commit, tag-N-gHASH on a commit after one, the short commit when no
# tag is reachable. A version that still comes out empty, dev or unknown
# fails the build. .git goes to /git, not /src/.git, where go build would
# find it and record VCS details of a work tree holding only backend/.
COPY .git /git
ARG VERSION
RUN version="${VERSION:-$(git --git-dir=/git describe --tags --always)}"; \
case "$version" in ""|dev|unknown) \
echo "version is '$version' although .git is present" >&2; \
exit 1 ;; \
esac; \
VERSION="$version" make build
# Frontend stage # Frontend stage
# node:22-alpine as of 2026-02-22 # node:22-alpine as of 2026-02-22
@@ -52,8 +65,9 @@ RUN yarn install --frozen-lockfile
RUN apk add --no-cache git make RUN apk add --no-cache git make
COPY . . COPY . .
# make frontend-check is the frontend half of make check (test + lint + # 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 # fmt-check); its test step runs the unit tests, then the production
# produces dist/ and gates the image on lint/fmt-check/test regressions. # 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 # This node stage has neither Go nor Docker; the lint and builder stages
# above gate the backend half. # above gate the backend half.
RUN make frontend-check RUN make frontend-check
+11 -5
View File
@@ -52,8 +52,8 @@ halves, so the root `make check` fails if either one is broken. We provide:
- `script/fmt` — format all files (writes): prettier, then gofmt over `backend/` - `script/fmt` — format all files (writes): prettier, then gofmt over `backend/`
- `script/fmt-check` — check formatting (read-only): prettier, then gofmt - `script/fmt-check` — check formatting (read-only): prettier, then gofmt
- `script/check` — run test, lint, and fmt-check - `script/check` — run test, lint, and fmt-check
- `script/frontend-test` — run the production build as the frontend's test (no - `script/frontend-test` — run the unit tests in `test/unit/` with Node's
unit tests yet) built-in test runner, then the production build
- `script/frontend-lint` — run prettier in check mode - `script/frontend-lint` — run prettier in check mode
- `script/frontend-fmt` — format everything prettier understands (writes) - `script/frontend-fmt` — format everything prettier understands (writes)
- `script/frontend-fmt-check` — check prettier formatting (read-only) - `script/frontend-fmt-check` — check prettier formatting (read-only)
@@ -136,8 +136,14 @@ Local hosts are tracked separately from WAN stats.
### Latency measurement ### Latency measurement
HEAD requests with `mode: 'no-cors'` and `cache: 'no-store'`, timed with HEAD requests with `mode: 'no-cors'` and `cache: 'no-store'`, timed with
`performance.now()`. 1-second timeout; anything over 1000ms is clamped to `performance.now()`. Each check times out after 80% of the refresh interval (24
unreachable. IPv4 only. 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 ### Color coding
@@ -212,7 +218,7 @@ What the [upaas](https://git.eeqj.de/sneak/upaas) app for netwatch needs:
- `REPORTS_PER_MINUTE`, default `60`: reports each client address may send a - `REPORTS_PER_MINUTE`, default `60`: reports each client address may send a
minute minute
- `DATA_DIR_MAX_BYTES`, default `1073741824` (1 GiB): the most room the - `DATA_DIR_MAX_BYTES`, default `1073741824` (1 GiB): the most room the
report files may take report files may take; the oldest are deleted to stay under it
- `CORS_ALLOWED_ORIGINS`, default empty: other origins whose pages may call - `CORS_ALLOWED_ORIGINS`, default empty: other origins whose pages may call
the API the API
- `DEBUG`, default `false`: debug logging - `DEBUG`, default `false`: debug logging
+17
View File
@@ -23,6 +23,23 @@ latest run passes.
# Completed Steps # Completed Steps
- 2026-10-03: `DATA_DIR_MAX_BYTES` is now how much of the report files is kept
(issue #54): when a report would take them past it, the oldest report files
are deleted to make room, each deletion logged, and at start files already
past it are deleted the same way. A file still being written is never deleted.
A report is refused with 507 only when the reports 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: 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): - 2026-09-29: the container sets up its own data directory (issue #75):
`bin/entrypoint.sh`, still as root, creates `DATA_DIR` if missing and gives it `bin/entrypoint.sh`, still as root, creates `DATA_DIR` if missing and gives it
and `/data` to the `netwatch` user with mode 750 before starting the backend and `/data` to the `netwatch` user with mode 750 before starting the backend
+13 -8
View File
@@ -32,7 +32,7 @@ pattern as the repo root: the targets in `backend/Makefile` are thin shims over
`test`, `fmt` and `fmt-check`: `test`, `fmt` and `fmt-check`:
- `script/build` — compile the static `netwatch-server` binary with its version - `script/build` — compile the static `netwatch-server` binary with its version
and architecture stamped in. The version is `VERSION` from the environment; stamped in. The version is `VERSION` from the environment;
when that is unset or empty, it falls back to `git describe` inside a git when that is unset or empty, it falls back to `git describe` inside a git
checkout, then to `dev` checkout, then to `dev`
- `script/test` — run the Go tests under a 30-second timeout - `script/test` — run the Go tests under a 30-second timeout
@@ -80,7 +80,7 @@ Internal packages in `internal/` follow standard Go project layout:
| `BIND_ADDRESS` | empty | IP address to listen on; empty listens on every interface | | `BIND_ADDRESS` | empty | IP address to listen on; empty listens on every interface |
| `PORT` | `8080` | HTTP listen port | | `PORT` | `8080` | HTTP listen port |
| `DATA_DIR` | `./data/reports` | Directory for compressed reports | | `DATA_DIR` | `./data/reports` | Directory for compressed reports |
| `DATA_DIR_MAX_BYTES` | `1073741824` (1 GiB) | Largest total size of the report files in `DATA_DIR`; see [Report limits](#report-limits) | | `DATA_DIR_MAX_BYTES` | `1073741824` (1 GiB) | Most bytes of report files kept in `DATA_DIR`, oldest deleted first; see [Report limits](#report-limits) |
| `DEBUG` | `false` | Enable debug logging | | `DEBUG` | `false` | Enable debug logging |
| `TRUSTED_PROXIES` | loopback + RFC1918 | Comma-separated CIDRs whose `X-Forwarded-For` / `X-Real-IP` headers are trusted for client IP resolution | | `TRUSTED_PROXIES` | loopback + RFC1918 | Comma-separated CIDRs whose `X-Forwarded-For` / `X-Real-IP` headers are trusted for client IP resolution |
| `REPORTS_PER_MINUTE` | `60` | Reports each client address may send a minute; see [Report limits](#report-limits) | | `REPORTS_PER_MINUTE` | `60` | Reports each client address may send a minute; see [Report limits](#report-limits) |
@@ -148,11 +148,17 @@ credentials, so it is bounded instead. Both refusals below answer with the same
`X-RateLimit-Reset` headers. `X-RateLimit-Reset` headers.
- **Size cap.** The report files in `DATA_DIR` may total at most - **Size cap.** The report files in `DATA_DIR` may total at most
`DATA_DIR_MAX_BYTES`, counting the files already there at start. Reports `DATA_DIR_MAX_BYTES`, counting the files already there at start. Reports
waiting in memory count at their uncompressed size until they are written, so waiting in memory count at their uncompressed size until they are written;
a report that would take the total past the cap is refused with 507, and those lost to a failed write stop counting, and the part of its file written
nothing of it is stored. Deleting report files frees room only at the next is removed. When a report would take the total past the cap, the oldest report
start, when the files are counted again. The default of 1 GiB is small enough files are deleted to make room, and each deletion is logged with the file's
for any host; set it to the space you can give `DATA_DIR`. name and size; a file still being written is never deleted. A report is
refused with 507, and nothing of it is stored, only when the reports waiting
to be written fill the cap on their own, and then no file is deleted. At
start, report files past the cap, as after lowering it, are deleted the same
way. So the cap is how much of the newest reports is kept: the default of 1
GiB is small enough for any host; set it to the space you can give
`DATA_DIR`.
### CORS ### CORS
@@ -170,7 +176,6 @@ starting, with an error naming `CORS_ALLOWED_ORIGINS`.
- Add integration test that POSTs a report and verifies the compressed output - Add integration test that POSTs a report and verifies the compressed output
- Add report decompression/query endpoint - Add report decompression/query endpoint
- Add metrics (Prometheus) for buffer size, flush count, report count - Add metrics (Prometheus) for buffer size, flush count, report count
- Add retention policy to prune old report files
## License ## License
-2
View File
@@ -21,7 +21,6 @@ import (
var ( var (
Appname = "netwatch-server" Appname = "netwatch-server"
Version string Version string
Buildarch string
) )
func main() { func main() {
@@ -40,7 +39,6 @@ func main() {
globals.Appname = Appname globals.Appname = Appname
globals.Version = Version globals.Version = Version
globals.Buildarch = Buildarch
fx.New( fx.New(
fx.Provide( fx.Provide(
-4
View File
@@ -10,22 +10,18 @@ var (
Appname string Appname string
// Version is the git version tag. // Version is the git version tag.
Version string Version string
// Buildarch is the build architecture.
Buildarch string
) )
// Globals holds build-time metadata for the application. // Globals holds build-time metadata for the application.
type Globals struct { type Globals struct {
Appname string Appname string
Version string Version string
Buildarch string
} }
// New creates a Globals instance from package-level variables. // New creates a Globals instance from package-level variables.
func New(_ fx.Lifecycle) (*Globals, error) { func New(_ fx.Lifecycle) (*Globals, error) {
return &Globals{ return &Globals{
Appname: Appname, Appname: Appname,
Buildarch: Buildarch,
Version: Version, Version: Version,
}, nil }, nil
} }
+4 -3
View File
@@ -86,11 +86,12 @@ func (s *Handlers) decodeErrorStatus(err error) int {
} }
// appendErrorStatus logs a failure to store a report and returns // appendErrorStatus logs a failure to store a report and returns
// the status to send: 507 when the report files are at their size // the status to send: 507 when the reports waiting to be written fill
// cap, otherwise 500. // the size cap, otherwise 500.
func (s *Handlers) appendErrorStatus(err error) int { func (s *Handlers) appendErrorStatus(err error) int {
if errors.Is(err, reportbuf.ErrFull) { if errors.Is(err, reportbuf.ErrFull) {
s.log.Warn("report refused: report files at their size cap") s.log.Warn("report refused: " +
"reports waiting to be written fill the size cap")
return http.StatusInsufficientStorage return http.StatusInsufficientStorage
} }
+3 -3
View File
@@ -68,9 +68,9 @@ func TestHandleReportStorageFailureIsNon2xx(t *testing.T) {
} }
} }
// TestHandleReportFullIs507 checks the answer when the report files // TestHandleReportFullIs507 checks the answer when the reports waiting
// are at their size cap: 507 and the usual error body, which tells // to be written fill the size cap: 507 and the usual error body, which
// the client nothing more. // tells the client nothing more.
func TestHandleReportFullIs507(t *testing.T) { func TestHandleReportFullIs507(t *testing.T) {
t.Parallel() t.Parallel()
+2 -1
View File
@@ -5,6 +5,7 @@ package logger
import ( import (
"log/slog" "log/slog"
"os" "os"
"runtime"
"sneak.berlin/go/netwatch/internal/globals" "sneak.berlin/go/netwatch/internal/globals"
@@ -95,6 +96,6 @@ func (l *Logger) Identify() {
l.log.Info("starting", l.log.Info("starting",
"appname", l.params.Globals.Appname, "appname", l.params.Globals.Appname,
"version", l.params.Globals.Version, "version", l.params.Globals.Version,
"buildarch", l.params.Globals.Buildarch, "arch", runtime.GOARCH,
) )
} }
+11 -1
View File
@@ -1,6 +1,9 @@
package reportbuf package reportbuf
import "time" import (
"os"
"time"
)
// Flush writes the buffered reports to a file now, as the periodic // Flush writes the buffered reports to a file now, as the periodic
// flush does, so tests need not wait a minute for it. // flush does, so tests need not wait a minute for it.
@@ -13,3 +16,10 @@ func (b *Buffer) Flush() error {
func (b *Buffer) StopClock(at time.Time) { func (b *Buffer) StopClock(at time.Time) {
b.now = func() time.Time { return at } b.now = func() time.Time { return at }
} }
// OnFileCreated makes the buffer call fn with each report file it
// writes from now on, once the file is created and before anything is
// written to it.
func (b *Buffer) OnFileCreated(fn func(f *os.File)) {
b.fileCreated = fn
}
+145 -45
View File
@@ -1,5 +1,6 @@
// Package reportbuf accumulates telemetry reports in memory // Package reportbuf accumulates telemetry reports in memory
// and periodically flushes them to zstd-compressed JSONL files. // and periodically flushes them to zstd-compressed JSONL files,
// deleting the oldest files to keep them under a size cap.
package reportbuf package reportbuf
import ( import (
@@ -37,9 +38,10 @@ const (
fileSuffix = ".jsonl.zst" fileSuffix = ".jsonl.zst"
) )
// ErrFull is returned by Append when storing the report would // ErrFull is returned by Append when the reports waiting to be
// take the report files past the configured maximum size. // written leave no room for the report under the configured maximum
var ErrFull = errors.New("report files at their size cap") // size, however many report files are deleted.
var ErrFull = errors.New("reports waiting to be written fill the size cap")
// Params defines the dependencies for Buffer. // Params defines the dependencies for Buffer.
type Params struct { type Params struct {
@@ -49,12 +51,29 @@ type Params struct {
Logger *logger.Logger Logger *logger.Logger
} }
// reportFile is a report file that may be deleted to make room, with
// the size it counts for in usedBytes.
type reportFile struct {
name string
size int64
}
// Buffer accumulates JSON lines in memory and flushes them // Buffer accumulates JSON lines in memory and flushes them
// to zstd-compressed files on disk. // to zstd-compressed files on disk.
type Buffer struct { type Buffer struct {
buf bytes.Buffer buf bytes.Buffer
dataDir string dataDir string
done chan struct{} done chan struct{}
// 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,
// oldest first: those in dataDir at start, then each one this
// buffer writes, once it is complete. A file still being written
// is not among them. filesBytes is their total size.
files []reportFile
filesBytes int64
log *slog.Logger log *slog.Logger
maxBytes int64 maxBytes int64
mu sync.Mutex mu sync.Mutex
@@ -66,8 +85,8 @@ type Buffer struct {
seq atomic.Uint64 seq atomic.Uint64
stopOnce sync.Once stopOnce sync.Once
// usedBytes is what Append checks against maxBytes: the size // usedBytes is what Append checks against maxBytes: the size
// of the report files in dataDir, plus the reports not yet // of the report files in dataDir, plus the reports waiting to
// written to one at their uncompressed size. // be written to one at their uncompressed size.
usedBytes int64 usedBytes int64
} }
@@ -85,6 +104,7 @@ func New(
b := &Buffer{ b := &Buffer{
dataDir: dir, dataDir: dir,
done: make(chan struct{}), done: make(chan struct{}),
fileCreated: func(*os.File) {},
log: params.Logger.Get(), log: params.Logger.Get(),
maxBytes: params.Config.DataDirMaxBytes, maxBytes: params.Config.DataDirMaxBytes,
now: time.Now, now: time.Now,
@@ -97,12 +117,28 @@ func New(
return fmt.Errorf("create data dir: %w", err) return fmt.Errorf("create data dir: %w", err)
} }
// Report files left by earlier runs count too. // Report files left by earlier runs count too, and are
b.usedBytes, err = reportFilesSize(b.dataDir) // the first to be deleted to make room.
files, err := reportFiles(b.dataDir)
if err != nil { if err != nil {
return err return err
} }
b.mu.Lock()
b.files = files
for _, f := range files {
b.filesBytes += f.size
}
b.usedBytes = b.filesBytes
// The files may be past the cap, if it was lowered since
// the last run.
b.deleteOldestFiles(0)
b.mu.Unlock()
go b.flushLoop() go b.flushLoop()
return nil return nil
@@ -128,9 +164,10 @@ func New(
} }
// Append marshals v as a single JSON line and appends it to // Append marshals v as a single JSON line and appends it to
// the buffer. It stores nothing and returns ErrFull if the line // the buffer. If the line would take usedBytes past maxBytes, the
// would take usedBytes past maxBytes. If the buffer reaches the // oldest report files are deleted to make room; it stores nothing
// size threshold, it is drained and written to disk // and returns ErrFull if that cannot make room. If the buffer
// reaches the size threshold, it is drained and written to disk
// asynchronously. // asynchronously.
func (b *Buffer) Append(v any) error { func (b *Buffer) Append(v any) error {
line, err := json.Marshal(v) line, err := json.Marshal(v)
@@ -142,6 +179,8 @@ func (b *Buffer) Append(v any) error {
b.mu.Lock() b.mu.Lock()
b.deleteOldestFiles(lineBytes)
if b.usedBytes+lineBytes > b.maxBytes { if b.usedBytes+lineBytes > b.maxBytes {
b.mu.Unlock() b.mu.Unlock()
@@ -171,6 +210,37 @@ func (b *Buffer) Append(v any) error {
return nil return nil
} }
// deleteOldestFiles deletes report files, oldest first, until n more
// bytes fit under maxBytes. It deletes none when the reports waiting
// to be written leave no room for n even with every file gone, since
// that would lose the files for nothing. The caller must hold b.mu.
func (b *Buffer) deleteOldestFiles(n int64) {
for b.usedBytes+n > b.maxBytes && len(b.files) > 0 {
if b.usedBytes-b.filesBytes+n > b.maxBytes {
return
}
f := b.files[0]
b.files = b.files[1:]
b.filesBytes -= f.size
// A file already gone, deleted by hand, has freed its room too.
err := os.Remove(filepath.Join(b.dataDir, f.name))
if err != nil && !errors.Is(err, fs.ErrNotExist) {
// The file is still there, so it still counts. It is
// not tried again until the next start.
b.log.Error("delete report file failed",
"file", f.name, "error", err)
continue
}
b.usedBytes -= f.size
b.log.Info("deleted report file to make room",
"file", f.name, "bytes", f.size)
}
}
// flushLoop runs a ticker that periodically flushes buffered // flushLoop runs a ticker that periodically flushes buffered
// data to disk until the done channel is closed. // data to disk until the done channel is closed.
func (b *Buffer) flushLoop() { func (b *Buffer) flushLoop() {
@@ -217,15 +287,41 @@ func (b *Buffer) drainBuf() []byte {
return data return data
} }
// writeFile creates a timestamped zstd-compressed JSONL file // writeFile writes data, reports drained from the buffer, to a new
// in the data directory. // timestamped zstd-compressed JSONL file in the data directory.
func (b *Buffer) writeFile(data []byte) error { func (b *Buffer) writeFile(data []byte) error {
// The timestamp comes first, so the names sort by time; the number // The timestamp comes first, so the names sort by time; the number
// after it tells apart files named in the same millisecond. // after it tells apart files named in the same millisecond.
ts := b.now().UTC().Format("2006-01-02T15-04-05.000Z") ts := b.now().UTC().Format("2006-01-02T15-04-05.000Z")
name := fmt.Sprintf("%s%s-%d%s", filePrefix, ts, b.seq.Add(1), fileSuffix) name := fmt.Sprintf("%s%s-%d%s", filePrefix, ts, b.seq.Add(1), fileSuffix)
path := filepath.Join(b.dataDir, name)
size, err := b.createFile(filepath.Join(b.dataDir, name), data)
// The reports no longer wait to be written, so they stop counting
// at their uncompressed size. If the write failed they are lost;
// otherwise they count as the file, which from here on may be
// deleted to make room.
b.mu.Lock()
defer b.mu.Unlock()
b.usedBytes -= int64(len(data))
if err != nil {
return err
}
b.usedBytes += size
b.files = append(b.files, reportFile{name: name, size: size})
b.filesBytes += size
return nil
}
// createFile creates the file at path holding data compressed with
// zstd, and returns its size. If the write fails once the file is
// created, it removes the file, so that a failed write leaves nothing
// behind to take room.
func (b *Buffer) createFile(path string, data []byte) (int64, error) {
// path is built from the operator-supplied dataDir plus a // path is built from the operator-supplied dataDir plus a
// generated timestamp and number, so it carries no external input. // generated timestamp and number, so it carries no external input.
f, err := os.OpenFile( //nolint:gosec // see comment above f, err := os.OpenFile( //nolint:gosec // see comment above
@@ -234,61 +330,65 @@ func (b *Buffer) writeFile(data []byte) error {
filePerms, filePerms,
) )
if err != nil { if err != nil {
return fmt.Errorf("create report file: %w", err) return 0, fmt.Errorf("create report file: %w", err)
} }
// Closes the file on the early returns below. The success b.fileCreated(f)
// path closes it explicitly to check the error; closing it
// a second time here is harmless.
defer func() { _ = f.Close() }()
size, err := writeCompressed(f, data)
if err != nil {
_ = f.Close()
return 0, errors.Join(err, os.Remove(path))
}
err = f.Close()
if err != nil {
err = fmt.Errorf("close report file: %w", err)
return 0, errors.Join(err, os.Remove(path))
}
return size, nil
}
// writeCompressed writes data to f compressed with zstd, and returns
// the size of f.
func writeCompressed(f *os.File, data []byte) (int64, error) {
enc, err := zstd.NewWriter(f) enc, err := zstd.NewWriter(f)
if err != nil { if err != nil {
return fmt.Errorf("create zstd encoder: %w", err) return 0, fmt.Errorf("create zstd encoder: %w", err)
} }
_, err = enc.Write(data) _, err = enc.Write(data)
if err != nil { if err != nil {
_ = enc.Close() _ = enc.Close()
return fmt.Errorf("write compressed data: %w", err) return 0, fmt.Errorf("write compressed data: %w", err)
} }
err = enc.Close() err = enc.Close()
if err != nil { if err != nil {
return fmt.Errorf("close zstd encoder: %w", err) return 0, fmt.Errorf("close zstd encoder: %w", err)
} }
info, err := f.Stat() info, err := f.Stat()
if err != nil { if err != nil {
return fmt.Errorf("stat report file: %w", err) return 0, fmt.Errorf("stat report file: %w", err)
} }
err = f.Close() return info.Size(), nil
if err != nil {
return fmt.Errorf("close report file: %w", err)
} }
// The reports counted at their uncompressed size while they // reportFiles returns the report files in dir, oldest first:
// waited; now they count as the file. After a failed write they // os.ReadDir sorts them by name, and the names sort by time.
// stay counted as they were, which errs toward refusing reports func reportFiles(dir string) ([]reportFile, error) {
// early rather than letting the files pass the cap.
b.mu.Lock()
b.usedBytes += info.Size() - int64(len(data))
b.mu.Unlock()
return nil
}
// reportFilesSize returns the total size of the report files in
// dir.
func reportFilesSize(dir string) (int64, error) {
entries, err := os.ReadDir(dir) entries, err := os.ReadDir(dir)
if err != nil { if err != nil {
return 0, fmt.Errorf("read data dir: %w", err) return nil, fmt.Errorf("read data dir: %w", err)
} }
var total int64 files := make([]reportFile, 0, len(entries))
for _, entry := range entries { for _, entry := range entries {
name := entry.Name() name := entry.Name()
@@ -299,11 +399,11 @@ func reportFilesSize(dir string) (int64, error) {
info, err := entry.Info() info, err := entry.Info()
if err != nil { if err != nil {
return 0, fmt.Errorf("stat report file: %w", err) return nil, fmt.Errorf("stat report file: %w", err)
} }
total += info.Size() files = append(files, reportFile{name: name, size: info.Size()})
} }
return total, nil return files, nil
} }
+343 -50
View File
@@ -139,41 +139,19 @@ func lineBytes(t *testing.T, report any) int {
return len(line) + 1 return len(line) + 1
} }
// TestAppendPastCapIsRefused fills the cap with a report not yet
// written. The next report is refused, and the report file already in
// DATA_DIR is kept: it is smaller than a report, so deleting it could
// not make room.
func TestAppendPastCapIsRefused(t *testing.T) { func TestAppendPastCapIsRefused(t *testing.T) {
report := map[string]string{"id": "cap"}
t.Setenv("DATA_DIR", t.TempDir())
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(lineBytes(t, report)))
buf := startBuffer(t)
err := buf.Append(report)
if err != nil {
t.Fatalf("report that fills the cap exactly: %v", err)
}
err = buf.Append(report)
if !errors.Is(err, reportbuf.ErrFull) {
t.Fatalf("report past the cap: error = %v, want ErrFull", err)
}
}
// 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
report := map[string]string{"id": "cap"} report := map[string]string{"id": "cap"}
dir := t.TempDir() dir := t.TempDir()
earlier := reportFilePath(dir, 1)
writeBytes(t, filepath.Join(dir, "reports-2026-01-01T00-00-00.000Z.jsonl.zst"), writeBytes(t, earlier, 1)
earlierBytes)
writeBytes(t, filepath.Join(dir, "notes.txt"), 10*earlierBytes)
t.Setenv("DATA_DIR", dir) t.Setenv("DATA_DIR", dir)
t.Setenv("DATA_DIR_MAX_BYTES", t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(1+lineBytes(t, report)))
strconv.Itoa(earlierBytes+lineBytes(t, report)))
buf := startBuffer(t) buf := startBuffer(t)
@@ -186,6 +164,96 @@ func TestCapCountsReportFilesAlreadyInDataDir(t *testing.T) {
if !errors.Is(err, reportbuf.ErrFull) { if !errors.Is(err, reportbuf.ErrFull) {
t.Fatalf("report past the cap: error = %v, want ErrFull", err) t.Fatalf("report past the cap: error = %v, want ErrFull", err)
} }
if !exists(t, earlier) {
t.Fatal("report file deleted, though that could not make room")
}
}
// 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": "oldest"}
dir := t.TempDir()
oldest := reportFilePath(dir, 1)
kept := []string{
reportFilePath(dir, 2),
reportFilePath(dir, 3),
filepath.Join(dir, "notes.txt"),
}
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(3*fileBytes+lineBytes(t, report)))
buf := startBuffer(t)
err := buf.Append(report)
if err != nil {
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 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)
}
}
} }
// TestWrittenReportsCountAtFileSize checks that once reports are // TestWrittenReportsCountAtFileSize checks that once reports are
@@ -195,8 +263,9 @@ func TestWrittenReportsCountAtFileSize(t *testing.T) {
// Repetitive, so its file is far smaller than its JSON. // Repetitive, so its file is far smaller than its JSON.
report := map[string]string{"id": strings.Repeat("a", 1000)} report := map[string]string{"id": strings.Repeat("a", 1000)}
size := lineBytes(t, report) size := lineBytes(t, report)
dir := t.TempDir()
t.Setenv("DATA_DIR", t.TempDir()) t.Setenv("DATA_DIR", dir)
// Room for the report twice over only if the first one counts // Room for the report twice over only if the first one counts
// at its file's size by the time the second arrives. // at its file's size by the time the second arrives.
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(2*size-1)) t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(2*size-1))
@@ -217,13 +286,22 @@ func TestWrittenReportsCountAtFileSize(t *testing.T) {
if err != nil { if err != nil {
t.Fatalf("second report, after the first was written: %v", err) t.Fatalf("second report, after the first was written: %v", err)
} }
if !hasReportFile(t, dir) {
t.Fatal("first report's file deleted, though the second fit beside it")
}
} }
// TestWrittenReportsKeepCounting writes one report file after another // TestWrittenFilesDeletedToMakeRoom writes one report file after
// under a small cap: each report must be taken while the files on disk // another under a small cap. Every report must be taken; files are
// leave room for it, and refused once they do not. // deleted only when the report would not fit beside them, and the
func TestWrittenReportsKeepCounting(t *testing.T) { // files kept leave room for it.
const maxBytes = 200 func TestWrittenFilesDeletedToMakeRoom(t *testing.T) {
const (
maxBytes = 200
// One file each, which take far more than maxBytes together.
reports = 50
)
report := map[string]string{"id": "written"} report := map[string]string{"id": "written"}
size := int64(lineBytes(t, report)) size := int64(lineBytes(t, report))
@@ -234,23 +312,24 @@ func TestWrittenReportsKeepCounting(t *testing.T) {
buf := startBuffer(t) buf := startBuffer(t)
// Every file takes at least a byte, so they fill the cap within for range reports {
// maxBytes rounds. before := reportFilesBytes(t, dir)
for range maxBytes {
used := reportFilesBytes(t, dir)
err := buf.Append(report) err := buf.Append(report)
if used+size > maxBytes {
if !errors.Is(err, reportbuf.ErrFull) {
t.Fatalf("with %d bytes of report files: error = %v, "+
"want ErrFull", used, err)
}
return
}
if err != nil { if err != nil {
t.Fatalf("with %d bytes of report files: %v", used, err) t.Fatalf("with %d bytes of report files: %v", before, 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() err = buf.Flush()
@@ -258,8 +337,199 @@ func TestWrittenReportsKeepCounting(t *testing.T) {
t.Fatalf("flush: %v", err) t.Fatalf("flush: %v", err)
} }
} }
}
t.Fatal("the report files never filled the cap") // TestFailedDeletionStillCounts makes deleting the oldest report file
// fail. It is still there, so it still takes room, and the next oldest
// is deleted in its place.
func TestFailedDeletionStillCounts(t *testing.T) {
const fileBytes = 100
report := map[string]string{"id": "stuck"}
dir := t.TempDir()
stuck := reportFilePath(dir, 1)
next := reportFilePath(dir, 2)
newest := reportFilePath(dir, 3)
for _, path := range []string{stuck, next, newest} {
writeBytes(t, path, fileBytes)
}
t.Setenv("DATA_DIR", dir)
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(3*fileBytes))
buf := startBuffer(t)
// A directory that is not empty cannot be deleted, even by root,
// which the tests run as in the backend image.
err := os.Remove(stuck)
if err != nil {
t.Fatalf("remove %s: %v", stuck, err)
}
err = os.Mkdir(stuck, 0o750)
if err != nil {
t.Fatalf("make directory %s: %v", stuck, err)
}
writeBytes(t, filepath.Join(stuck, "file"), 1)
err = buf.Append(report)
if err != nil {
t.Fatalf("report past the cap: %v", err)
}
if exists(t, next) {
t.Fatal("next oldest report file kept: the failed deletion " +
"counted as making room")
}
if !exists(t, newest) {
t.Fatal("newest report file deleted, though deleting one made room")
}
}
// TestFileDeletedByHandFreesRoom deletes the oldest report file by
// hand after start. When room is needed, its room counts as freed, so
// no other file is deleted.
func TestFileDeletedByHandFreesRoom(t *testing.T) {
const fileBytes = 100
report := map[string]string{"id": "by-hand"}
dir := t.TempDir()
gone := reportFilePath(dir, 1)
kept := reportFilePath(dir, 2)
writeBytes(t, gone, fileBytes)
writeBytes(t, kept, fileBytes)
t.Setenv("DATA_DIR", dir)
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(2*fileBytes))
buf := startBuffer(t)
err := os.Remove(gone)
if err != nil {
t.Fatalf("remove %s: %v", gone, err)
}
err = buf.Append(report)
if err != nil {
t.Fatalf("report past the cap: %v", err)
}
if !exists(t, kept) {
t.Fatal("report file deleted, though the one deleted by hand " +
"had made room")
}
}
// TestFileBeingWrittenIsNeverDeleted holds a write open, its file
// created but not complete, while a report needs room. The file may
// not be deleted to make it, so the report is refused. Once the write
// is complete, the file is deleted when room is needed.
func TestFileBeingWrittenIsNeverDeleted(t *testing.T) {
report := map[string]string{"id": "writing"}
t.Setenv("DATA_DIR", t.TempDir())
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(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("report that fills the cap exactly: %v", err)
}
flushed := make(chan error)
go func() { flushed <- buf.Flush() }()
writing := <-created
// Errorf, not Fatalf, until the write is released, so that a
// failure here does not leave it held.
err = buf.Append(report)
if !errors.Is(err, reportbuf.ErrFull) {
t.Errorf("report while the first was being written: "+
"error = %v, want ErrFull", err)
}
if !exists(t, writing) {
t.Error("report file deleted while it was being written")
}
close(release)
err = <-flushed
if err != nil {
t.Fatalf("flush: %v", err)
}
// Writes from here on, the final one at stop included, go
// through unheld.
buf.OnFileCreated(func(*os.File) {})
err = buf.Append(report)
if err != nil {
t.Fatalf("report after the write was complete: %v", err)
}
if exists(t, writing) {
t.Fatal("complete report file kept when room was needed")
}
}
// TestFailedWriteStopsCounting makes a write fail once its file is
// created. Its reports are lost, so they stop counting, and the part of
// the file written is removed, so it takes no room.
func TestFailedWriteStopsCounting(t *testing.T) {
report := map[string]string{"id": "failed"}
t.Setenv("DATA_DIR", t.TempDir())
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(lineBytes(t, report)))
buf := startBuffer(t)
var failed string
// Closing the file under the write makes the write fail.
buf.OnFileCreated(func(f *os.File) {
failed = f.Name()
_ = f.Close()
})
err := buf.Append(report)
if err != nil {
t.Fatalf("report that fills the cap exactly: %v", err)
}
err = buf.Flush()
if err == nil {
t.Fatal("flush succeeded, though its file was closed under it")
}
if exists(t, failed) {
t.Fatal("file of the failed write kept")
}
// Writes from here on, the final one at stop included, succeed.
buf.OnFileCreated(func(*os.File) {})
err = buf.Append(report)
if err != nil {
t.Fatalf("report after the failed write: %v", err)
}
} }
// TestConcurrentAppendsStopAtCap appends from many goroutines at once // TestConcurrentAppendsStopAtCap appends from many goroutines at once
@@ -407,6 +677,14 @@ func readReportFiles(t *testing.T, dir string) []string {
return contents return contents
} }
// reportFilePath returns the path in dir of a report file named as
// written on the given day of January 2026, so that a lower day sorts
// as older.
func reportFilePath(dir string, day int) string {
return filepath.Join(dir,
fmt.Sprintf("reports-2026-01-%02dT00-00-00.000Z-1.jsonl.zst", day))
}
func writeBytes(t *testing.T, path string, n int) { func writeBytes(t *testing.T, path string, n int) {
t.Helper() t.Helper()
@@ -416,6 +694,21 @@ func writeBytes(t *testing.T, path string, n int) {
} }
} }
func exists(t *testing.T, path string) bool {
t.Helper()
_, err := os.Stat(path)
if errors.Is(err, fs.ErrNotExist) {
return false
}
if err != nil {
t.Fatalf("stat %s: %v", path, err)
}
return true
}
func hasReportFile(t *testing.T, dir string) bool { func hasReportFile(t *testing.T, dir string) bool {
t.Helper() t.Helper()
+2 -1
View File
@@ -4,6 +4,7 @@ import (
"errors" "errors"
"net" "net"
"net/http" "net/http"
"runtime"
"strconv" "strconv"
"time" "time"
@@ -53,7 +54,7 @@ func (s *Server) listenAndServe() {
s.log.Info("http begin listen", s.log.Info("http begin listen",
"listenaddr", s.httpServer.Addr, "listenaddr", s.httpServer.Addr,
"version", s.params.Globals.Version, "version", s.params.Globals.Version,
"buildarch", s.params.Globals.Buildarch, "arch", runtime.GOARCH,
) )
err := s.httpServer.ListenAndServe() err := s.httpServer.ListenAndServe()
+2 -2
View File
@@ -1,6 +1,6 @@
#!/bin/sh #!/bin/sh
# script/build: compile the static netwatch-server binary into the # script/build: compile the static netwatch-server binary into the
# backend project root, with its version and architecture stamped in. # backend project root, with its version stamped in.
set -eu set -eu
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)" ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
@@ -14,7 +14,7 @@ main() {
version="${VERSION:-$(git describe --always --dirty 2>/dev/null || echo dev)}" version="${VERSION:-$(git describe --always --dirty 2>/dev/null || echo dev)}"
CGO_ENABLED=0 go build -trimpath \ CGO_ENABLED=0 go build -trimpath \
-ldflags "-s -w -X main.Version=$version -X main.Buildarch=$(uname -m)" \ -ldflags "-s -w -X main.Version=$version" \
-o netwatch-server ./cmd/netwatch-server/ -o netwatch-server ./cmd/netwatch-server/
} }
+3
View File
@@ -9,6 +9,9 @@
type="image/svg+xml" 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>" 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> </head>
<body class="bg-gray-900 text-white min-h-screen"> <body class="bg-gray-900 text-white min-h-screen">
<div id="app"></div> <div id="app"></div>
+2 -3
View File
@@ -16,9 +16,8 @@ main() {
"$SCRIPT_DIR/check" "$SCRIPT_DIR/check"
# Own line: a failing command substitution inside an argument does # Own line: a failing command substitution inside an argument does
# not trip `set -e`, so the inline form degrades silently to an # not trip `set -e`, so the inline form degrades silently to an
# empty constant. VERSION is computed here because .dockerignore # empty constant. The VERSION build argument takes precedence over
# excludes .git, so `git describe` in a build stage yields an empty # the version a build stage derives from the .git in the context.
# version without failing.
version="$(git describe --tags --always --dirty 2>/dev/null || true)" version="$(git describe --tags --always --dirty 2>/dev/null || true)"
[ -n "$version" ] || version="unknown" [ -n "$version" ] || version="unknown"
docker build --no-cache \ docker build --no-cache \
+2 -3
View File
@@ -12,9 +12,8 @@ main() {
cd "$ROOT" cd "$ROOT"
# Own line: a failing command substitution inside an argument does # Own line: a failing command substitution inside an argument does
# not trip `set -e`, so the inline form degrades silently to an # not trip `set -e`, so the inline form degrades silently to an
# empty constant. VERSION is computed here because .dockerignore # empty constant. The VERSION build argument takes precedence over
# excludes .git, so `git describe` in a build stage yields an empty # the version a build stage derives from the .git in the context.
# version without failing.
version="$(git describe --tags --always --dirty 2>/dev/null || true)" version="$(git describe --tags --always --dirty 2>/dev/null || true)"
[ -n "$version" ] || version="unknown" [ -n "$version" ] || version="unknown"
docker build --no-cache \ docker build --no-cache \
+4 -3
View File
@@ -1,13 +1,14 @@
#!/bin/sh #!/bin/sh
# script/frontend-test: run the frontend test suite. The frontend has no # script/frontend-test: run the frontend test suite: the unit tests in
# unit tests; the production build serves as the test (fails on broken # test/unit/ with Node's built-in test runner, then the production
# code). # build, which fails on broken code.
set -eu set -eu
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)" ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
main() { main() {
cd "$ROOT" cd "$ROOT"
timeout 30 node --test test/unit/*.test.js
timeout 30 yarn build timeout 30 yarn build
} }
+51 -26
View File
@@ -1,14 +1,14 @@
import "./styles.css";
// --- Configuration ----------------------------------------------------------- // --- Configuration -----------------------------------------------------------
// Timing, axis labels, and display constants. Latency above maxLatency is // Timing, axis labels, and display constants. A target check times out
// clamped to "unreachable". The sparkline Y-axis is capped at // 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 // 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 // display their real value in the latency figure. The history buffer holds
// maxHistoryPoints samples (historyDuration / updateInterval). // maxHistoryPoints samples (historyDuration / updateInterval).
// reportInterval is how often collected samples are POSTed to the backend. // reportInterval is how often collected samples are POSTed to the backend.
const CONFIG = { export const CONFIG = {
updateInterval: 3000, updateInterval: 3000,
maxHistoryPoints: 100, maxHistoryPoints: 100,
reportInterval: 60000, reportInterval: 60000,
@@ -16,7 +16,7 @@ const CONFIG = {
return (this.maxHistoryPoints * this.updateInterval) / 1000; return (this.maxHistoryPoints * this.updateInterval) / 1000;
}, },
get requestTimeout() { get requestTimeout() {
return Math.min(this.updateInterval - 100, 3000); return this.updateInterval * 0.8;
}, },
get maxLatency() { get maxLatency() {
return this.requestTimeout; return this.requestTimeout;
@@ -503,12 +503,16 @@ class Reporter {
// --- Latency Measurement ----------------------------------------------------- // --- 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 controller = new AbortController();
const timeoutId = setTimeout( const timeoutId = setTimeout(
() => controller.abort(), () => controller.abort(),
CONFIG.requestTimeout, CONFIG.requestTimeout,
); );
signal?.addEventListener("abort", () => controller.abort());
const targetUrl = new URL(url); const targetUrl = new URL(url);
targetUrl.searchParams.set("_cb", Date.now().toString()); targetUrl.searchParams.set("_cb", Date.now().toString());
@@ -1101,7 +1105,7 @@ function sortAndRebuildWAN(state) {
// --- Main Loop --------------------------------------------------------------- // --- Main Loop ---------------------------------------------------------------
async function tick(state, onOffline) { async function tick(state, signal, onOffline) {
const ts = Date.now(); const ts = Date.now();
if (state.paused) { if (state.paused) {
@@ -1123,11 +1127,12 @@ async function tick(state, onOffline) {
log.debug(`Tick #${state.tickCount + 1} started`); log.debug(`Tick #${state.tickCount + 1} started`);
const results = await Promise.all( 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 // User may have paused, or the next round may have given up this
if (state.paused) return; // one's checks, while awaiting results — discard them
if (state.paused || signal.aborted) return;
state.tickCount++; state.tickCount++;
@@ -1168,9 +1173,10 @@ async function tick(state, onOffline) {
// --- Recovery Probe ---------------------------------------------------------- // --- Recovery Probe ----------------------------------------------------------
// When offline, rapidly poll 4 random WAN hosts every 500ms. As soon as any // When offline, check 4 random WAN hosts every 500ms, giving up the checks
// responds, stop probing and fire a normal tick to refresh all hosts. // started 500ms before, so at most 4 are ever waiting. As soon as one
function startRecoveryProbe(state, triggerTick) { // answers, stop probing and start a new round at once.
function startRecoveryProbe(state, startRounds) {
if (state._recoveryProbeId) return; // already running if (state._recoveryProbeId) return; // already running
const candidates = [...state.wan]; const candidates = [...state.wan];
for (let i = candidates.length - 1; i > 0; i--) { for (let i = candidates.length - 1; i > 0; i--) {
@@ -1181,15 +1187,18 @@ function startRecoveryProbe(state, triggerTick) {
log.notice( log.notice(
`Recovery probe started (${canaries.map((h) => h.name).join(", ")})`, `Recovery probe started (${canaries.map((h) => h.name).join(", ")})`,
); );
state._recoveryProbeId = setInterval(async () => { state._recoveryProbeId = setInterval(() => {
if (state.paused) return; if (state.paused) return;
const results = await Promise.all( state._recoveryProbeChecks?.abort();
canaries.map((h) => measureLatency(h.url)), const checks = new AbortController();
); state._recoveryProbeChecks = checks;
if (results.some((r) => r.error === null)) { 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"); log.notice("Recovery probe: connectivity detected");
stopRecoveryProbe(state); stopRecoveryProbe(state);
triggerTick(); startRounds();
});
} }
}, 500); }, 500);
} }
@@ -1198,6 +1207,7 @@ function stopRecoveryProbe(state) {
if (state._recoveryProbeId) { if (state._recoveryProbeId) {
clearInterval(state._recoveryProbeId); clearInterval(state._recoveryProbeId);
state._recoveryProbeId = null; state._recoveryProbeId = null;
state._recoveryProbeChecks?.abort();
} }
} }
@@ -1374,18 +1384,34 @@ async function init() {
updateClocks(); updateClocks();
setInterval(updateClocks, 1000); 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() { function doTick() {
tick(state, () => startRecoveryProbe(state, doTick)); roundChecks.abort();
roundChecks = new AbortController();
tick(state, roundChecks.signal, () =>
startRecoveryProbe(state, startRounds),
);
} }
// Starts a round now and then one every CONFIG.updateInterval.
let tickIntervalId;
function startRounds() {
clearInterval(tickIntervalId);
doTick(); doTick();
let tickIntervalId = setInterval(doTick, CONFIG.updateInterval); tickIntervalId = setInterval(doTick, CONFIG.updateInterval);
}
startRounds();
document document
.getElementById("interval-select") .getElementById("interval-select")
.addEventListener("change", (e) => { .addEventListener("change", (e) => {
const newInterval = parseInt(e.target.value, 10); const newInterval = parseInt(e.target.value, 10);
clearInterval(tickIntervalId);
CONFIG.updateInterval = newInterval; CONFIG.updateInterval = newInterval;
log.notice( log.notice(
`Interval changed to ${humanDuration(newInterval / 1000)}, history reset`, `Interval changed to ${humanDuration(newInterval / 1000)}, history reset`,
@@ -1434,8 +1460,7 @@ async function init() {
// Start immediately with new interval // Start immediately with new interval
stopRecoveryProbe(state); stopRecoveryProbe(state);
doTick(); startRounds();
tickIntervalId = setInterval(doTick, CONFIG.updateInterval);
}); });
window.addEventListener("resize", () => handleResize(state)); 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 // 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 // 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 (typeof document !== "undefined" && document.getElementById("app")) {
if (document.readyState === "loading") { if (document.readyState === "loading") {
document.addEventListener("DOMContentLoaded", init); document.addEventListener("DOMContentLoaded", init);
+66
View File
@@ -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",
});
});
}