2 Commits
Author SHA1 Message Date
clawbot 304cf79826 fix(backend): give each report file a name of its own (closes #61)
check / check (push) Successful in 55s
Report files were named by a millisecond timestamp and created with
O_EXCL, so two flushes in the same millisecond, such as a flush for
size and the final flush at shutdown, got the same name and the second
failed, losing its reports. Each name now carries a number after the
timestamp that goes up by one for each file the server starts to
write, so names still sort by time and never repeat within a run; a
failed write leaves a gap in the numbers. The buffer reads the time
through a clock the new test stops, so its two flushes share one
timestamp on every run; the 1 ms pauses earlier tests used to dodge
the collision are gone.

Model: opus-5-5
2026-09-29 05:06:38 +00:00
clawbot d2f219ca19 upaas: health check, settings checked at start, README section (closes #59)
check / check (push) Successful in 15s
The image's HEALTHCHECK requests /.well-known/healthcheck through
nginx on the port from PORT, so it fails unless both processes answer.
The backend reads PORT and DEBUG with strconv instead of viper, which
turned a bad PORT into 0 and a bad DEBUG into false. Those, and a
BIND_ADDRESS that is not an IP address, now stop the start with an
error naming the variable; the TRUSTED_PROXIES error names it too.
bin/entrypoint.sh also refuses a container PORT outside 1 to 65535,
or 8081, where the backend listens, naming PORT. README.md gains
"Running under upaas". Its first-run steps create the host directory
owned by uid 1000, so the image changes no ownership.

Model: opus-5-5
2026-09-29 06:39:10 +02:00
12 changed files with 301 additions and 38 deletions
+6
View File
@@ -83,6 +83,12 @@ VOLUME /data
# The default public port; PORT changes it.
EXPOSE 8080
# Requests the backend's health check through nginx, on the port from
# PORT, so it fails unless both answer. upaas reads the result 60
# seconds after a deploy and fails the deploy unless it is healthy.
HEALTHCHECK --interval=30s --timeout=5s --start-period=10s --retries=3 \
CMD wget -q -O /dev/null "http://127.0.0.1:${PORT:-8080}/.well-known/healthcheck"
# The nginx image stops its container with SIGQUIT; the entrypoint
# acts on TERM and INT.
STOPSIGNAL SIGTERM
+40
View File
@@ -194,6 +194,46 @@ only inside the container, on `127.0.0.1:8081`. The image:
- Writes buffered reports to disk on `docker stop`, and exits non-zero if nginx
or the backend exits on its own, so the platform restarts it
## Running under upaas
What the [upaas](https://git.eeqj.de/sneak/upaas) app for netwatch needs:
- **Port:** container port `8080`.
- **Volume:** container path `/data`; the reports are kept in `/data/reports`.
- **First run:** upaas bind-mounts the host directory it is given and does not
create it, and the backend, which runs as uid 1000, does not start unless it
can write there. Create the directory, owned by uid 1000, before the first
deploy:
```bash
mkdir -p /path/to/data
chown 1000:1000 /path/to/data
```
- **Environment variables:** none is required. An empty one counts as unset, and
one set to a value netwatch cannot use stops the container at start, with the
reason in its log.
- `PORT`, default `8080`: the container port, from 1 to 65535. `8081` cannot
be used: the backend listens on it inside the container
- `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
- `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
- `TRUSTED_PROXIES`, default loopback and RFC1918: leave unset. The
backend's only client is nginx, on loopback, which passes on the client
address; nginx takes it from `X-Forwarded-For` only from RFC1918
addresses.
- **Health check:** the image's `HEALTHCHECK` requests
`/.well-known/healthcheck` through nginx every 30 seconds, so it fails unless
both nginx and the backend answer. upaas reads the container's health 60
seconds after a deploy and fails the deploy unless it is `healthy`. The
container also stops when either process exits.
## Browser Compatibility
Requires a modern browser with ES modules, Fetch API, Canvas API, and CSS custom
+16
View File
@@ -23,6 +23,22 @@ latest run passes.
# Completed Steps
- 2026-09-29: report file names can no longer collide (issue #61): each is
`reports-<timestamp>-<number>.jsonl.zst`, where the number goes up by one for
each file the server starts to write, so two flushes in the same millisecond,
such as a flush for size and the final flush at shutdown, each get a file of
their own instead of the second one failing. A failed write leaves a gap in
the numbers.
- 2026-09-29: ready to run under upaas (issue #59): the image has a
`HEALTHCHECK` that requests `/.well-known/healthcheck` through nginx on the
port from `PORT`. The backend no longer reads a bad `PORT` as 0 or a bad
`DEBUG` as false: those, and a `BIND_ADDRESS` that is not an IP address, stop
it from starting with an error naming the variable, as the limits,
`CORS_ALLOWED_ORIGINS` and, now by name, `TRUSTED_PROXIES` already did.
`bin/entrypoint.sh` also refuses a `PORT` outside 1 to 65535, and `8081`,
where the backend listens inside the container, naming `PORT`. `README.md` has
a "Running under upaas" section, whose first-run steps create the host
directory for `/data` owned by uid 1000; the image does not change its owner
- 2026-09-29: nginx listens on `PORT` (issue #26), 8080 when unset or empty: the
nginx image renders `nginx.conf` as a template at container start, filling in
`PORT` and no other variable. `bin/entrypoint.sh` refuses to start when `PORT`
+11 -3
View File
@@ -91,6 +91,10 @@ The loopback entries cover the reverse proxy that shares the container; the
RFC1918 ranges match `nginx.conf`. A request whose direct peer is outside this
set has its forwarded headers ignored, and the direct peer is logged instead.
A variable set to a value the server cannot use, such as `PORT=abc`,
`DEBUG=maybe` or a `BIND_ADDRESS` that is not an IP address, stops it from
starting, with an error naming the variable. An empty variable counts as unset.
### Container image
The root `Dockerfile` builds one image in which nginx listens on the public port
@@ -102,9 +106,13 @@ which `netwatch` owns.
### Report storage
Reports are written as `reports-<timestamp>.jsonl.zst` files in `DATA_DIR`.
Each file contains one JSON object per line, compressed with zstd. Files are
created with `O_EXCL` to prevent overwrites.
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, and a failed write leaves a gap in the numbers. Each
file contains one JSON object per line, compressed with zstd. Files are created
with `O_EXCL` to prevent overwrites.
### Report limits
+51 -16
View File
@@ -6,7 +6,10 @@ import (
"errors"
"fmt"
"log/slog"
"math"
"net/netip"
"net/url"
"strconv"
"strings"
"sneak.berlin/go/netwatch/internal/globals"
@@ -37,6 +40,9 @@ var (
errNotOrigin = errors.New(
"must be an origin, scheme://host with an optional port",
)
errNotPort = errors.New("must be a port number, 1 to 65535")
errNotBool = errors.New("must be true or false")
errNotIP = errors.New("must be an IP address, or empty")
)
// Params defines the dependencies for Config.
@@ -65,7 +71,8 @@ type Config struct {
}
// New loads configuration from env, .env files, and config
// files, returning a fully resolved Config.
// files, returning a fully resolved Config. It fails, with an error
// naming the setting, on a value the server cannot use.
func New(
_ fx.Lifecycle,
params Params,
@@ -103,15 +110,29 @@ func New(
}
}
// Read with strconv: viper's GetInt and GetBool would read a value
// they cannot parse as 0 or false instead of failing.
port, err := strconv.Atoi(viper.GetString("PORT"))
if err != nil || port < 1 || port > math.MaxUint16 {
return nil, fmt.Errorf("PORT %q: %w",
viper.GetString("PORT"), errNotPort)
}
debug, err := strconv.ParseBool(viper.GetString("DEBUG"))
if err != nil {
return nil, fmt.Errorf("DEBUG %q: %w",
viper.GetString("DEBUG"), errNotBool)
}
s := &Config{
BindAddress: viper.GetString("BIND_ADDRESS"),
CORSAllowedOrigins: splitList(viper.GetString("CORS_ALLOWED_ORIGINS")),
DataDir: viper.GetString("DATA_DIR"),
DataDirMaxBytes: viper.GetInt64("DATA_DIR_MAX_BYTES"),
Debug: viper.GetBool("DEBUG"),
Debug: debug,
MetricsPassword: viper.GetString("METRICS_PASSWORD"),
MetricsUsername: viper.GetString("METRICS_USERNAME"),
Port: viper.GetInt("PORT"),
Port: port,
ReportsPerMinute: viper.GetInt("REPORTS_PER_MINUTE"),
SentryDSN: viper.GetString("SENTRY_DSN"),
TrustedProxies: splitList(viper.GetString("TRUSTED_PROXIES")),
@@ -119,19 +140,7 @@ func New(
params: &params,
}
// viper reads a value that is not a number as 0, so this also
// catches a mistyped setting.
if s.ReportsPerMinute <= 0 {
return nil, fmt.Errorf("REPORTS_PER_MINUTE %q: %w",
viper.GetString("REPORTS_PER_MINUTE"), errNotPositive)
}
if s.DataDirMaxBytes <= 0 {
return nil, fmt.Errorf("DATA_DIR_MAX_BYTES %q: %w",
viper.GetString("DATA_DIR_MAX_BYTES"), errNotPositive)
}
err = checkOrigins(s.CORSAllowedOrigins)
err = s.check()
if err != nil {
return nil, err
}
@@ -144,6 +153,32 @@ func New(
return s, nil
}
// check fails with an error naming the first setting here whose value
// the server cannot use. New checks PORT and DEBUG as it reads them,
// and the middleware checks TRUSTED_PROXIES as it parses it.
func (s *Config) check() error {
// viper reads a value that is not a number as 0, so this also
// catches a mistyped setting.
if s.ReportsPerMinute <= 0 {
return fmt.Errorf("REPORTS_PER_MINUTE %q: %w",
viper.GetString("REPORTS_PER_MINUTE"), errNotPositive)
}
if s.DataDirMaxBytes <= 0 {
return fmt.Errorf("DATA_DIR_MAX_BYTES %q: %w",
viper.GetString("DATA_DIR_MAX_BYTES"), errNotPositive)
}
if s.BindAddress != "" {
_, err := netip.ParseAddr(s.BindAddress)
if err != nil {
return fmt.Errorf("BIND_ADDRESS %q: %w", s.BindAddress, errNotIP)
}
}
return checkOrigins(s.CORSAllowedOrigins)
}
// checkOrigins fails on the first CORS_ALLOWED_ORIGINS entry that is
// not a plain origin, scheme://host with an optional port, as browsers
// send it; anything more, such as a trailing "/", would match no page.
+57
View File
@@ -29,6 +29,63 @@ func requireConfigError(t *testing.T, setting string) {
}
}
// TestSettingsLoadAsGiven: valid values pass the checks and are used
// as given. bin/entrypoint.sh starts the server with these
// BIND_ADDRESS and PORT values.
func TestSettingsLoadAsGiven(t *testing.T) {
t.Setenv("BIND_ADDRESS", "127.0.0.1")
t.Setenv("PORT", "8081")
t.Setenv("DEBUG", "true")
var cfg *config.Config
app := fx.New(
fx.NopLogger,
fx.Provide(globals.New, logger.New, config.New),
fx.Populate(&cfg),
)
err := app.Err()
if err != nil {
t.Fatalf("config error = %v", err)
}
if cfg.BindAddress != "127.0.0.1" || cfg.Port != 8081 || !cfg.Debug {
t.Fatalf("BindAddress, Port, Debug = %q, %d, %t; "+
"want \"127.0.0.1\", 8081, true",
cfg.BindAddress, cfg.Port, cfg.Debug)
}
}
// TestPortMustBeAPortNumber: viper reads a value that is not a number
// as 0, on which the server would listen on a random port.
func TestPortMustBeAPortNumber(t *testing.T) {
for _, value := range []string{"abc", "0", "65536", "8080.5"} {
t.Run(value, func(t *testing.T) {
t.Setenv("PORT", value)
requireConfigError(t, "PORT")
})
}
}
// TestDebugMustBeTrueOrFalse: viper reads any other value, such as
// "yes", as false.
func TestDebugMustBeTrueOrFalse(t *testing.T) {
t.Setenv("DEBUG", "yes")
requireConfigError(t, "DEBUG")
}
// TestBindAddressMustBeAnIPAddress: a host name would be looked up
// only once the server starts listening, and a mistyped one would stop
// it then with an error that does not name the setting.
func TestBindAddressMustBeAnIPAddress(t *testing.T) {
t.Setenv("BIND_ADDRESS", "localhost")
requireConfigError(t, "BIND_ADDRESS")
}
// TestReportsPerMinuteMustBePositive: unchecked, zero would panic
// when the routes are built, and a negative rate would lift the
// limit.
+3 -3
View File
@@ -77,8 +77,8 @@ func New(
return s, nil
}
// parseTrustedProxies converts CIDR strings into prefixes,
// failing fast on any malformed entry.
// parseTrustedProxies converts the TRUSTED_PROXIES entries into
// prefixes, failing fast on any malformed entry.
func parseTrustedProxies(cidrs []string) ([]netip.Prefix, error) {
prefixes := make([]netip.Prefix, 0, len(cidrs))
@@ -86,7 +86,7 @@ func parseTrustedProxies(cidrs []string) ([]netip.Prefix, error) {
prefix, err := netip.ParsePrefix(cidr)
if err != nil {
return nil, fmt.Errorf(
"trusted proxy %q: %w", cidr, err,
"TRUSTED_PROXIES %q: %w", cidr, err,
)
}
@@ -40,8 +40,8 @@ func TestParseTrustedProxiesRejectsMalformed(t *testing.T) {
t.Parallel()
_, err := middleware.ParseTrustedProxies([]string{"not-a-cidr"})
if err == nil {
t.Fatal("expected error for malformed CIDR, got nil")
if err == nil || !strings.Contains(err.Error(), "TRUSTED_PROXIES") {
t.Fatalf("error = %v, want one naming TRUSTED_PROXIES", err)
}
}
@@ -1,7 +1,15 @@
package reportbuf
import "time"
// Flush writes the buffered reports to a file now, as the periodic
// flush does, so tests need not wait a minute for it.
func (b *Buffer) Flush() error {
return b.flushLocked()
}
// StopClock makes every report file the buffer writes from now on
// carry the timestamp at, as if all were written in one millisecond.
func (b *Buffer) StopClock(at time.Time) {
b.now = func() time.Time { return at }
}
+16 -4
View File
@@ -14,6 +14,7 @@ import (
"path/filepath"
"strings"
"sync"
"sync/atomic"
"time"
"sneak.berlin/go/netwatch/internal/config"
@@ -30,7 +31,8 @@ const (
dirPerms fs.FileMode = 0o750
filePerms fs.FileMode = 0o640
// Report files are named filePrefix + timestamp + fileSuffix.
// Report files are named filePrefix + timestamp + "-" + number +
// fileSuffix; see writeFile.
filePrefix = "reports-"
fileSuffix = ".jsonl.zst"
)
@@ -56,6 +58,12 @@ type Buffer struct {
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
// seq numbers the report files, so that two named in the same
// millisecond still get different names.
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
@@ -79,6 +87,7 @@ func New(
done: make(chan struct{}),
log: params.Logger.Get(),
maxBytes: params.Config.DataDirMaxBytes,
now: time.Now,
}
lc.Append(fx.Hook{
@@ -211,11 +220,14 @@ func (b *Buffer) drainBuf() []byte {
// writeFile creates a timestamped zstd-compressed JSONL file
// in the data directory.
func (b *Buffer) writeFile(data []byte) error {
ts := time.Now().UTC().Format("2006-01-02T15-04-05.000Z")
path := filepath.Join(b.dataDir, filePrefix+ts+fileSuffix)
// 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)
// path is built from the operator-supplied dataDir plus a
// generated timestamp, 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
path,
os.O_WRONLY|os.O_CREATE|os.O_EXCL,
+77 -8
View File
@@ -3,9 +3,11 @@ package reportbuf_test
import (
"encoding/json"
"errors"
"fmt"
"io/fs"
"os"
"path/filepath"
"slices"
"strconv"
"strings"
"sync"
@@ -18,6 +20,7 @@ import (
"sneak.berlin/go/netwatch/internal/logger"
"sneak.berlin/go/netwatch/internal/reportbuf"
"github.com/klauspost/compress/zstd"
"go.uber.org/fx"
"go.uber.org/fx/fxtest"
)
@@ -210,10 +213,6 @@ func TestWrittenReportsCountAtFileSize(t *testing.T) {
t.Fatalf("flush: %v", err)
}
// The second report is written at shutdown, and must not land in
// the first file's millisecond (see TestWrittenReportsKeepCounting).
time.Sleep(time.Millisecond)
err = buf.Append(report)
if err != nil {
t.Fatalf("second report, after the first was written: %v", err)
@@ -254,10 +253,6 @@ func TestWrittenReportsKeepCounting(t *testing.T) {
t.Fatalf("with %d bytes of report files: %v", used, err)
}
// Report files are named to the millisecond; two in the same
// one collide (https://git.eeqj.de/sneak/netwatch/issues/61).
time.Sleep(time.Millisecond)
err = buf.Flush()
if err != nil {
t.Fatalf("flush: %v", err)
@@ -315,6 +310,43 @@ func TestConcurrentAppendsStopAtCap(t *testing.T) {
}
}
// TestTwoFlushesInOneMillisecond flushes twice within one millisecond,
// as a flush for size and the final flush at shutdown can: each flush
// must write a file of its own, and the files must hold every report.
func TestTwoFlushesInOneMillisecond(t *testing.T) {
const flushes = 2
dir := t.TempDir()
t.Setenv("DATA_DIR", dir)
buf := startBuffer(t)
buf.StopClock(time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC))
for id := 1; id <= flushes; id++ {
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 %d: %v", id, err)
}
}
files := readReportFiles(t, dir)
if len(files) != flushes {
t.Fatalf("%d report files after %d flushes", len(files), flushes)
}
for id := 1; id <= flushes; id++ {
want := fmt.Sprintf(`{"id":%d}`+"\n", id)
if !slices.Contains(files, want) {
t.Fatalf("no report file holds report %d alone", id)
}
}
}
// reportFilesBytes returns the total size of the report files in dir.
func reportFilesBytes(t *testing.T, dir string) int64 {
t.Helper()
@@ -338,6 +370,43 @@ func reportFilesBytes(t *testing.T, dir string) int64 {
return total
}
// readReportFiles returns the decompressed contents of each report
// file in dir.
func readReportFiles(t *testing.T, dir string) []string {
t.Helper()
files := os.DirFS(dir)
names, err := fs.Glob(files, "reports-*.jsonl.zst")
if err != nil {
t.Fatalf("list report files: %v", err)
}
dec, err := zstd.NewReader(nil)
if err != nil {
t.Fatalf("create zstd decoder: %v", err)
}
defer dec.Close()
contents := make([]string, 0, len(names))
for _, name := range names {
compressed, readErr := fs.ReadFile(files, name)
if readErr != nil {
t.Fatalf("read %s: %v", name, readErr)
}
data, decErr := dec.DecodeAll(compressed, nil)
if decErr != nil {
t.Fatalf("decompress %s: %v", name, decErr)
}
contents = append(contents, string(data))
}
return contents
}
func writeBytes(t *testing.T, path string, n int) {
t.Helper()
+14 -2
View File
@@ -10,8 +10,9 @@ set -u
# PORT is the public port nginx listens on, 8080 when unset or empty.
# nginx would take a value such as localhost or unix:/tmp/x.sock as an
# address and start anyway, so anything but digits stops the container
# here, before either process starts.
# address and start anyway, and reports a bad port without naming
# PORT, so a value that is not a usable port stops the container here,
# before either process starts.
export PORT="${PORT:-8080}"
case "$PORT" in
*[!0-9]*)
@@ -19,6 +20,17 @@ case "$PORT" in
exit 1
;;
esac
# The length is checked first because, for a number too big for it,
# the shell's test prints an error and is false, so the range checks
# alone would let it through.
if [ "${#PORT}" -gt 5 ] || [ "$PORT" -lt 1 ] || [ "$PORT" -gt 65535 ]; then
echo "entrypoint: PORT must be from 1 to 65535, not '$PORT'" >&2
exit 1
fi
if [ "$PORT" -eq 8081 ]; then
echo "entrypoint: PORT cannot be 8081, netwatch-server listens there" >&2
exit 1
fi
# A stop signal is only noted here; the loop below acts on it.
stop_requested=""