6 Commits
Author SHA1 Message Date
clawbot 8833603eff nginx: trust X-Forwarded-For only from TRUSTED_PROXIES (closes #64)
check / check (push) Successful in 15s
nginx trusted X-Forwarded-For from every RFC1918 address, so a client
reaching it from one could write a new address on each request and
get a fresh rate-limit allowance. The container's TRUSTED_PROXIES now
names the reverse proxies nginx trusts, none by default.
bin/entrypoint.sh makes each entry a CIDR, checks it with the new
"netwatch-server check-cidr", which runs the server's own
TRUSTED_PROXIES parsing, and writes one set_real_ip_from line per
entry into /etc/nginx/trusted-proxies.conf, which nginx.conf includes.
The backend is started with TRUSTED_PROXIES=127.0.0.1/32, since nginx
is its only client. The viewport test mounts an empty file there.

Model: opus-5-5
2026-09-29 08:55:47 +02:00
clawbot 6022cc8b02 fix(backend): give each report file a name of its own (closes #61)
check / check (push) Successful in 13s
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 uses up its number, leaving a gap if the file could not
be created and otherwise a file under that number that may be
incomplete.

Model: opus-5-5
2026-09-29 08:05:26 +02: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
clawbot ced1956b06 nginx: listen on PORT, default 8080; server_tokens off (closes #26)
check / check (push) Successful in 15s
nginx.conf is now a template the nginx image renders into conf.d at
container start. bin/entrypoint.sh sets PORT to 8080 when unset or
empty, and stops with an error before starting anything when PORT is
not digits only: nginx would take a value such as localhost or
unix:/tmp/x.sock as an address and start anyway. NGINX_ENVSUBST_FILTER
limits the rendering to PORT, so $uri, $host and every other nginx
variable pass through unchanged. server_tokens off drops the version
from the Server header and error pages. script/frontend-viewport-test
renders the template the same way. EXPOSE still documents 8080; the
backend stays on 127.0.0.1:8081.

Model: opus-5-5
2026-09-29 04:55:48 +02:00
clawbot ea66caf338 fix(backend): rate-limit and cap report ingest, drop wildcard CORS (closes #20)
check / check (push) Successful in 11s
POST /api/v1/reports stays unauthenticated but is bounded. Each client
address, as the trusted-proxy logic resolves it, may send
REPORTS_PER_MINUTE reports a minute (default 60, counted by
go-chi/httprate over a sliding minute); past that it gets 429 with
Retry-After. reportbuf refuses a report that would take the report
files past DATA_DIR_MAX_BYTES (default 1 GiB), counting the files
already in DATA_DIR and unwritten reports at their uncompressed size;
the handler answers 507. CORS adds nothing unless CORS_ALLOWED_ORIGINS
lists origins. A limit that is not a positive number, or an origin
that is not a plain scheme://host[:port], stops the server from
starting.

Model: opus-5-5
2026-09-29 04:22:19 +02:00
clawbot bbcc7d921d build: one image, nginx in front of the backend on loopback (closes #52)
check / check (push) Successful in 12s
The root Dockerfile builds the only image; Dockerfile.backend is gone.
Its stages: lint, a Go stage that runs the tests and builds
netwatch-server, the node stage, and an nginx runtime. nginx serves
dist/ on 8080 and proxies /api/ and /.well-known/healthcheck to the
backend on 127.0.0.1:8081. bin/entrypoint.sh starts both, turns TERM or
INT into a stop of both, and exits non-zero when either exits on its
own. The backend runs as user netwatch and keeps reports on the /data
volume. New setting BIND_ADDRESS (empty: every interface). STOPSIGNAL is
SIGTERM, since the nginx image's SIGQUIT would miss the entrypoint.
script/docker is the org model verbatim.

Model: opus-5-5
2026-09-29 02:59:33 +02:00
22 changed files with 1346 additions and 97 deletions
+12 -3
View File
@@ -63,13 +63,15 @@ RUN make frontend-check
FROM nginx@sha256:15e96e59aa3b0aada3a121296e3bce117721f42d88f5f64217ef4b18f458c6ab FROM nginx@sha256:15e96e59aa3b0aada3a121296e3bce117721f42d88f5f64217ef4b18f458c6ab
# netwatch-server runs as this user, which owns the report directory. # netwatch-server runs as this user, which owns the report directory.
# nginx keeps the image's own arrangement: master process as root, # nginx keeps the image's own arrangement: its main process runs as
# workers as the nginx user. # root, its worker processes as the nginx user.
RUN addgroup -g 1000 -S netwatch && \ RUN addgroup -g 1000 -S netwatch && \
adduser -u 1000 -S netwatch -G netwatch adduser -u 1000 -S netwatch -G netwatch
# At start-up the nginx image renders every template here into
# conf.d; bin/entrypoint.sh says how.
RUN rm /etc/nginx/conf.d/default.conf RUN rm /etc/nginx/conf.d/default.conf
COPY nginx.conf /etc/nginx/conf.d/netwatch.conf COPY nginx.conf /etc/nginx/templates/netwatch.conf.template
COPY --from=frontend /app/dist /usr/share/nginx/html COPY --from=frontend /app/dist /usr/share/nginx/html
COPY --from=builder /src/netwatch-server /usr/local/bin/netwatch-server COPY --from=builder /src/netwatch-server /usr/local/bin/netwatch-server
COPY bin/entrypoint.sh /usr/local/bin/entrypoint.sh COPY bin/entrypoint.sh /usr/local/bin/entrypoint.sh
@@ -78,8 +80,15 @@ ENV DATA_DIR=/data/reports
RUN mkdir -p /data/reports && chown -R netwatch:netwatch /data RUN mkdir -p /data/reports && chown -R netwatch:netwatch /data
VOLUME /data VOLUME /data
# The default public port; PORT changes it.
EXPOSE 8080 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 # The nginx image stops its container with SIGQUIT; the entrypoint
# acts on TERM and INT. # acts on TERM and INT.
STOPSIGNAL SIGTERM STOPSIGNAL SIGTERM
+48 -2
View File
@@ -184,8 +184,8 @@ container: nginx serves the built frontend and passes `/api/` and
only inside the container, on `127.0.0.1:8081`. The image: only inside the container, on `127.0.0.1:8081`. The image:
- Listens on port 8080 by default (override with `PORT` env var) - Listens on port 8080 by default (override with `PORT` env var)
- Trusts `X-Forwarded-For` from RFC1918 reverse proxies (10/8, 172.16/12, - Takes the client address from `X-Forwarded-For` only on requests from the
192.168/16) reverse proxies named in `TRUSTED_PROXIES`, and by default from none
- Sends access logs to stdout - Sends access logs to stdout
- Caches static assets with immutable headers - Caches static assets with immutable headers
- Stores reports in `DATA_DIR`, `/data/reports` by default, on the `/data` - Stores reports in `DATA_DIR`, `/data/reports` by default, on the `/data`
@@ -194,6 +194,52 @@ 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 - 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 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 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
`X-Forwarded-For` only on a request from one of them, and the rate limit
counts that address. Unset, `X-Forwarded-For` is ignored and every client
behind the proxy shares the proxy's one allowance of `REPORTS_PER_MINUTE`.
Name only addresses nothing but the proxy connects from: any client that
connects from one can write its own `X-Forwarded-For`, and through a port
Docker publishes, every client may connect from the Docker network's
gateway, such as `172.17.0.1`.
- **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 ## Browser Compatibility
Requires a modern browser with ES modules, Fetch API, Canvas API, and CSS custom Requires a modern browser with ES modules, Fetch API, Canvas API, and CSS custom
+42
View File
@@ -23,6 +23,48 @@ latest run passes.
# Completed Steps # Completed Steps
- 2026-09-29: nginx takes the client address from `X-Forwarded-For` only on
requests from the reverse proxies named in the container's `TRUSTED_PROXIES`
(issue #64), and by default from none, where it trusted every RFC1918 address
before, so a client could write a new address on each request and escape the
rate limit. `bin/entrypoint.sh` writes one `set_real_ip_from` line per entry
into `/etc/nginx/trusted-proxies.conf`, which `nginx.conf` includes, refusing
an entry that is not an IP address or CIDR, as `netwatch-server check-cidr`
finds; it starts the backend with `TRUSTED_PROXIES=127.0.0.1/32`, since nginx
is its only client
- 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 uses up its
number, leaving a gap if the file could not be created and otherwise a file
under that number that may be incomplete.
- 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`
is not digits only. `server_tokens off` keeps the nginx version out of
responses. `script/frontend-viewport-test` renders the template the same way.
Gzip and a `50x.html` error page are not added
- 2026-09-29: bounded the report endpoint (issue #20): `POST /api/v1/reports`
still needs no credentials, but each client address, as resolved through
`TRUSTED_PROXIES`, may send `REPORTS_PER_MINUTE` (default 60) reports a
minute, counted by `go-chi/httprate`, and past that gets 429 with
`Retry-After`; the report files in `DATA_DIR`, counted from start with those
already there, may total at most `DATA_DIR_MAX_BYTES` (default 1 GiB), past
which reports get 507; and the wildcard CORS is gone: no CORS headers unless
`CORS_ALLOWED_ORIGINS` lists origins, and an entry that is not a plain
`scheme://host[:port]` origin, `*` included, stops the server from starting.
Deleting report files frees room only at the next start; pruning is issue #54
- 2026-09-28: one container image (issue #52): the root `Dockerfile` builds the - 2026-09-28: one container image (issue #52): the root `Dockerfile` builds the
only image, and `Dockerfile.backend` is gone. nginx serves the frontend on only image, and `Dockerfile.backend` is gone. nginx serves the frontend on
port 8080 and proxies `/api/` and `/.well-known/healthcheck` to the backend, port 8080 and proxies `/api/` and `/.well-known/healthcheck` to the backend,
+74 -15
View File
@@ -75,18 +75,26 @@ Internal packages in `internal/` follow standard Go project layout:
### Configuration ### Configuration
| Variable | Default | Description | | Variable | Default | Description |
| ----------------- | -------------------- | -------------------------------------------------------------------------------------------------------- | | ---------------------- | -------------------- | -------------------------------------------------------------------------------------------------------- |
| `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 |
| `DEBUG` | `false` | Enable debug logging | | `DATA_DIR_MAX_BYTES` | `1073741824` (1 GiB) | Largest total size of the report files in `DATA_DIR`; see [Report limits](#report-limits) |
| `TRUSTED_PROXIES` | loopback + RFC1918 | Comma-separated CIDRs whose `X-Forwarded-For` / `X-Real-IP` headers are trusted for client IP resolution | | `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) |
| `CORS_ALLOWED_ORIGINS` | empty | Comma-separated origins whose pages may call the API; see [CORS](#cors) |
`TRUSTED_PROXIES` defaults to `127.0.0.1/32,::1/128,10.0.0.0/8,172.16.0.0/12,192.168.0.0/16`. `TRUSTED_PROXIES` defaults to `127.0.0.1/32,::1/128,10.0.0.0/8,172.16.0.0/12,192.168.0.0/16`.
The loopback entries cover the reverse proxy that shares the container; the The loopback entries cover a reverse proxy on the same host. A request whose
RFC1918 ranges match `nginx.conf`. A request whose direct peer is outside this direct peer is outside this set has its forwarded headers ignored, and the
set has its forwarded headers ignored, and the direct peer is logged instead. direct peer is logged and rate-limited instead. The container image does not use
this default; see [Container image](#container-image).
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 ### Container image
@@ -94,14 +102,65 @@ The root `Dockerfile` builds one image in which nginx listens on the public port
8080, serves the frontend, and proxies `/api/` and `/.well-known/healthcheck` to 8080, serves the frontend, and proxies `/api/` and `/.well-known/healthcheck` to
this server. The image's entrypoint, `bin/entrypoint.sh`, starts the server as this server. The image's entrypoint, `bin/entrypoint.sh`, starts the server as
user `netwatch` (uid 1000) with `BIND_ADDRESS=127.0.0.1` and `PORT=8081`, so user `netwatch` (uid 1000) with `BIND_ADDRESS=127.0.0.1` and `PORT=8081`, so
only nginx reaches it. `DATA_DIR` is `/data/reports`, on the `/data` volume, only nginx reaches it, and with `TRUSTED_PROXIES=127.0.0.1/32`, so it takes the
which `netwatch` owns. client address nginx passes on and no other. `DATA_DIR` is `/data/reports`, on
the `/data` volume, which `netwatch` owns.
The container's own `TRUSTED_PROXIES` goes to nginx instead: IP addresses or
CIDRs, separated by commas, of the reverse proxies in front of the container.
nginx takes the client address from `X-Forwarded-For` only on a request from one
of them. Unset or empty, nginx trusts no proxy, and the client address is the
one each request comes from, so every client behind a proxy shares one rate
limit. An entry that is not an IP address or CIDR, such as a hostname or
`1.2.3`, stops the container at start with an error naming `TRUSTED_PROXIES`:
the entrypoint checks each entry with `netwatch-server check-cidr`, which parses
it as this server parses its own `TRUSTED_PROXIES`.
### Report storage ### Report storage
Reports are written as `reports-<timestamp>.jsonl.zst` files in `DATA_DIR`. Reports are written as `reports-<timestamp>-<number>.jsonl.zst` files in
Each file contains one JSON object per line, compressed with zstd. Files are `DATA_DIR`. The timestamp is in UTC to the millisecond, so the names sort by
created with `O_EXCL` to prevent overwrites. 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.
### Report limits
`POST /api/v1/reports` takes reports from anyone who can reach it, without
credentials, so it is bounded instead. Both refusals below answer with the same
`{"status":"error"}` body as any other error.
- **Rate limit.** Each client address, resolved through `TRUSTED_PROXIES`, may
send `REPORTS_PER_MINUTE` reports a minute; past that it gets 429 with
`Retry-After: 60`. The minute slides: reports from the minute before still
count, fading out over the current one, so an address is sure never to be
refused only while it sends at most half of `REPORTS_PER_MINUTE` in any 60
seconds. The page sends one report a minute from each open tab, so the default
of 60 refuses nothing from up to 30 tabs behind one address, such as a
household or an office sharing it, however their reports bunch up. Report
responses also carry `X-RateLimit-Limit`, `X-RateLimit-Remaining` and
`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`.
### CORS
The page calls the API from the origin it is served from, so by default the
server sends no CORS headers, and browsers let no other origin's pages call it.
To serve the page from elsewhere, list that origin in `CORS_ALLOWED_ORIGINS`
(for example `https://netwatch.example.com`); pages from a listed origin may
`GET` and `POST` with a `Content-Type` header. Each entry must be a plain
origin, `scheme://host` with an optional `:port`, as browsers send it: no path,
not even a trailing `/`, and no `*`. Any other entry stops the server from
starting, with an error naming `CORS_ALLOWED_ORIGINS`.
## TODO ## TODO
+16
View File
@@ -2,6 +2,9 @@
package main package main
import ( import (
"fmt"
"os"
"sneak.berlin/go/netwatch/internal/config" "sneak.berlin/go/netwatch/internal/config"
"sneak.berlin/go/netwatch/internal/globals" "sneak.berlin/go/netwatch/internal/globals"
"sneak.berlin/go/netwatch/internal/handlers" "sneak.berlin/go/netwatch/internal/handlers"
@@ -22,6 +25,19 @@ var (
) )
func main() { func main() {
// "netwatch-server check-cidr CIDR" exits 1, with the error, if
// this server would refuse CIDR in its TRUSTED_PROXIES.
// bin/entrypoint.sh runs it on each entry it gives nginx.
if len(os.Args) == 3 && os.Args[1] == "check-cidr" {
_, err := middleware.ParseTrustedProxies(os.Args[2:])
if err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
return
}
globals.Appname = Appname globals.Appname = Appname
globals.Version = Version globals.Version = Version
globals.Buildarch = Buildarch globals.Buildarch = Buildarch
+4 -1
View File
@@ -5,6 +5,7 @@ go 1.25.5
require ( require (
github.com/go-chi/chi/v5 v5.2.5 github.com/go-chi/chi/v5 v5.2.5
github.com/go-chi/cors v1.2.2 github.com/go-chi/cors v1.2.2
github.com/go-chi/httprate v0.16.0
github.com/joho/godotenv v1.5.1 github.com/joho/godotenv v1.5.1
github.com/klauspost/compress v1.18.4 github.com/klauspost/compress v1.18.4
github.com/spf13/viper v1.21.0 github.com/spf13/viper v1.21.0
@@ -14,6 +15,7 @@ require (
require ( require (
github.com/fsnotify/fsnotify v1.9.0 // indirect github.com/fsnotify/fsnotify v1.9.0 // indirect
github.com/go-viper/mapstructure/v2 v2.4.0 // indirect github.com/go-viper/mapstructure/v2 v2.4.0 // indirect
github.com/klauspost/cpuid/v2 v2.2.10 // indirect
github.com/pelletier/go-toml/v2 v2.2.4 // indirect github.com/pelletier/go-toml/v2 v2.2.4 // indirect
github.com/sagikazarmark/locafero v0.11.0 // indirect github.com/sagikazarmark/locafero v0.11.0 // indirect
github.com/sourcegraph/conc v0.3.1-0.20240121214520-5f936abd7ae8 // indirect github.com/sourcegraph/conc v0.3.1-0.20240121214520-5f936abd7ae8 // indirect
@@ -21,10 +23,11 @@ require (
github.com/spf13/cast v1.10.0 // indirect github.com/spf13/cast v1.10.0 // indirect
github.com/spf13/pflag v1.0.10 // indirect github.com/spf13/pflag v1.0.10 // indirect
github.com/subosito/gotenv v1.6.0 // indirect github.com/subosito/gotenv v1.6.0 // indirect
github.com/zeebo/xxh3 v1.0.2 // indirect
go.uber.org/dig v1.19.0 // indirect go.uber.org/dig v1.19.0 // indirect
go.uber.org/multierr v1.10.0 // indirect go.uber.org/multierr v1.10.0 // indirect
go.uber.org/zap v1.26.0 // indirect go.uber.org/zap v1.26.0 // indirect
go.yaml.in/yaml/v3 v3.0.4 // indirect go.yaml.in/yaml/v3 v3.0.4 // indirect
golang.org/x/sys v0.29.0 // indirect golang.org/x/sys v0.30.0 // indirect
golang.org/x/text v0.28.0 // indirect golang.org/x/text v0.28.0 // indirect
) )
+10 -2
View File
@@ -8,6 +8,8 @@ github.com/go-chi/chi/v5 v5.2.5 h1:Eg4myHZBjyvJmAFjFvWgrqDTXFyOzjj7YIm3L3mu6Ug=
github.com/go-chi/chi/v5 v5.2.5/go.mod h1:X7Gx4mteadT3eDOMTsXzmI4/rwUpOwBHLpAfupzFJP0= github.com/go-chi/chi/v5 v5.2.5/go.mod h1:X7Gx4mteadT3eDOMTsXzmI4/rwUpOwBHLpAfupzFJP0=
github.com/go-chi/cors v1.2.2 h1:Jmey33TE+b+rB7fT8MUy1u0I4L+NARQlK6LhzKPSyQE= github.com/go-chi/cors v1.2.2 h1:Jmey33TE+b+rB7fT8MUy1u0I4L+NARQlK6LhzKPSyQE=
github.com/go-chi/cors v1.2.2/go.mod h1:sSbTewc+6wYHBBCW7ytsFSn836hqM7JxpglAy2Vzc58= github.com/go-chi/cors v1.2.2/go.mod h1:sSbTewc+6wYHBBCW7ytsFSn836hqM7JxpglAy2Vzc58=
github.com/go-chi/httprate v0.16.0 h1:8V5DH9j6pSK6UQoBsTpvMyFxycqaKEIToyPKzHJjUa8=
github.com/go-chi/httprate v0.16.0/go.mod h1:A8lo+qRhk+s9LiuP5saS7XCGDXRXMcrueq0NfIuCa/I=
github.com/go-viper/mapstructure/v2 v2.4.0 h1:EBsztssimR/CONLSZZ04E8qAkxNYq4Qp9LvH92wZUgs= github.com/go-viper/mapstructure/v2 v2.4.0 h1:EBsztssimR/CONLSZZ04E8qAkxNYq4Qp9LvH92wZUgs=
github.com/go-viper/mapstructure/v2 v2.4.0/go.mod h1:oJDH3BJKyqBA2TXFhDsKDGDTlndYOZ6rGS0BRZIxGhM= github.com/go-viper/mapstructure/v2 v2.4.0/go.mod h1:oJDH3BJKyqBA2TXFhDsKDGDTlndYOZ6rGS0BRZIxGhM=
github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI=
@@ -16,6 +18,8 @@ github.com/joho/godotenv v1.5.1 h1:7eLL/+HRGLY0ldzfGMeQkb7vMd0as4CfYvUVzLqw0N0=
github.com/joho/godotenv v1.5.1/go.mod h1:f4LDr5Voq0i2e/R5DDNOoa2zzDfwtkZa6DnEwAbqwq4= github.com/joho/godotenv v1.5.1/go.mod h1:f4LDr5Voq0i2e/R5DDNOoa2zzDfwtkZa6DnEwAbqwq4=
github.com/klauspost/compress v1.18.4 h1:RPhnKRAQ4Fh8zU2FY/6ZFDwTVTxgJ/EMydqSTzE9a2c= github.com/klauspost/compress v1.18.4 h1:RPhnKRAQ4Fh8zU2FY/6ZFDwTVTxgJ/EMydqSTzE9a2c=
github.com/klauspost/compress v1.18.4/go.mod h1:R0h/fSBs8DE4ENlcrlib3PsXS61voFxhIs2DeRhCvJ4= github.com/klauspost/compress v1.18.4/go.mod h1:R0h/fSBs8DE4ENlcrlib3PsXS61voFxhIs2DeRhCvJ4=
github.com/klauspost/cpuid/v2 v2.2.10 h1:tBs3QSyvjDyFTq3uoc/9xFpCuOsJQFNPiAhYdw2skhE=
github.com/klauspost/cpuid/v2 v2.2.10/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0=
github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE=
github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk=
github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
@@ -42,6 +46,10 @@ github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/subosito/gotenv v1.6.0 h1:9NlTDc1FTs4qu0DDq7AEtTPNw6SVm7uBMsUCUjABIf8= github.com/subosito/gotenv v1.6.0 h1:9NlTDc1FTs4qu0DDq7AEtTPNw6SVm7uBMsUCUjABIf8=
github.com/subosito/gotenv v1.6.0/go.mod h1:Dk4QP5c2W3ibzajGcXpNraDfq2IrhjMIvMSWPKKo0FU= github.com/subosito/gotenv v1.6.0/go.mod h1:Dk4QP5c2W3ibzajGcXpNraDfq2IrhjMIvMSWPKKo0FU=
github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ=
github.com/zeebo/assert v1.3.0/go.mod h1:Pq9JiuJQpG8JLJdtkwrJESF0Foym2/D9XMU5ciN/wJ0=
github.com/zeebo/xxh3 v1.0.2 h1:xZmwmqxHZA8AI603jOQ0tMqmBr9lPeFwGg6d+xy9DC0=
github.com/zeebo/xxh3 v1.0.2/go.mod h1:5NWz9Sef7zIDm2JHfFlcQvNekmcEl9ekUZQQKCYaDcA=
go.uber.org/dig v1.19.0 h1:BACLhebsYdpQ7IROQ1AGPjrXcP5dF80U3gKoFzbaq/4= go.uber.org/dig v1.19.0 h1:BACLhebsYdpQ7IROQ1AGPjrXcP5dF80U3gKoFzbaq/4=
go.uber.org/dig v1.19.0/go.mod h1:Us0rSJiThwCv2GteUN0Q7OKvU7n5J4dxZ9JKUXozFdE= go.uber.org/dig v1.19.0/go.mod h1:Us0rSJiThwCv2GteUN0Q7OKvU7n5J4dxZ9JKUXozFdE=
go.uber.org/fx v1.24.0 h1:wE8mruvpg2kiiL1Vqd0CC+tr0/24XIB10Iwp2lLWzkg= go.uber.org/fx v1.24.0 h1:wE8mruvpg2kiiL1Vqd0CC+tr0/24XIB10Iwp2lLWzkg=
@@ -54,8 +62,8 @@ go.uber.org/zap v1.26.0 h1:sI7k6L95XOKS281NhVKOFCUNIvv9e0w4BF8N3u+tCRo=
go.uber.org/zap v1.26.0/go.mod h1:dtElttAiwGvoJ/vj4IwHBS/gXsEu/pZ50mUIRWuG0so= go.uber.org/zap v1.26.0/go.mod h1:dtElttAiwGvoJ/vj4IwHBS/gXsEu/pZ50mUIRWuG0so=
go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc=
go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg=
golang.org/x/sys v0.29.0 h1:TPYlXGxvx1MGTn2GiZDhnjPA9wZzZeGKHHmKhHYvgaU= golang.org/x/sys v0.30.0 h1:QjkSwP/36a20jFYWkSue1YwXzLmsV5Gfq7Eiy72C1uc=
golang.org/x/sys v0.29.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/sys v0.30.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/text v0.28.0 h1:rhazDwis8INMIwQ4tpjLDzUhx6RlXqZNPEM0huQojng= golang.org/x/text v0.28.0 h1:rhazDwis8INMIwQ4tpjLDzUhx6RlXqZNPEM0huQojng=
golang.org/x/text v0.28.0/go.mod h1:U8nCwOR8jO/marOQ0QbDiOngZVEBB7MAiitBuMjXiNU= golang.org/x/text v0.28.0/go.mod h1:U8nCwOR8jO/marOQ0QbDiOngZVEBB7MAiitBuMjXiNU=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
+123 -25
View File
@@ -4,7 +4,12 @@ package config
import ( import (
"errors" "errors"
"fmt"
"log/slog" "log/slog"
"math"
"net/netip"
"net/url"
"strconv"
"strings" "strings"
"sneak.berlin/go/netwatch/internal/globals" "sneak.berlin/go/netwatch/internal/globals"
@@ -16,13 +21,31 @@ import (
) )
// defaultTrustedProxies lists the networks whose forwarded // defaultTrustedProxies lists the networks whose forwarded
// headers are honoured by default. It covers the RFC1918 // headers are honoured by default: IPv4 and IPv6 loopback,
// ranges (to match nginx.conf) plus IPv4 and IPv6 loopback, // for a reverse proxy on the same host, and the RFC1918
// because the reverse proxy shares the container and reaches // ranges. The container image does not use it:
// the backend over loopback. // bin/entrypoint.sh gives the server 127.0.0.1/32, since
// nginx is its only client there.
const defaultTrustedProxies = "127.0.0.1/32,::1/128," + const defaultTrustedProxies = "127.0.0.1/32,::1/128," +
"10.0.0.0/8,172.16.0.0/12,192.168.0.0/16" "10.0.0.0/8,172.16.0.0/12,192.168.0.0/16"
// Default limits on stored reports; backend/README.md gives the
// reasons for these values.
const (
defaultReportsPerMinute = 60
defaultDataDirMaxBytes = 1 << 30 // 1 GiB
)
var (
errNotPositive = errors.New("must be a positive whole number")
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. // Params defines the dependencies for Config.
type Params struct { type Params struct {
fx.In fx.In
@@ -33,20 +56,24 @@ type Params struct {
// Config holds the resolved application configuration. // Config holds the resolved application configuration.
type Config struct { type Config struct {
BindAddress string BindAddress string
DataDir string CORSAllowedOrigins []string
Debug bool DataDir string
MetricsPassword string DataDirMaxBytes int64
MetricsUsername string Debug bool
Port int MetricsPassword string
SentryDSN string MetricsUsername string
TrustedProxies []string Port int
log *slog.Logger ReportsPerMinute int
params *Params SentryDSN string
TrustedProxies []string
log *slog.Logger
params *Params
} }
// New loads configuration from env, .env files, and config // 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( func New(
_ fx.Lifecycle, _ fx.Lifecycle,
params Params, params Params,
@@ -61,11 +88,15 @@ func New(
viper.AutomaticEnv() viper.AutomaticEnv()
// An empty CORS_ALLOWED_ORIGINS allows no other origin.
viper.SetDefault("CORS_ALLOWED_ORIGINS", "")
viper.SetDefault("DATA_DIR", "./data/reports") viper.SetDefault("DATA_DIR", "./data/reports")
viper.SetDefault("DATA_DIR_MAX_BYTES", defaultDataDirMaxBytes)
viper.SetDefault("DEBUG", "false") viper.SetDefault("DEBUG", "false")
// An empty BIND_ADDRESS listens on every interface. // An empty BIND_ADDRESS listens on every interface.
viper.SetDefault("BIND_ADDRESS", "") viper.SetDefault("BIND_ADDRESS", "")
viper.SetDefault("PORT", "8080") viper.SetDefault("PORT", "8080")
viper.SetDefault("REPORTS_PER_MINUTE", defaultReportsPerMinute)
viper.SetDefault("SENTRY_DSN", "") viper.SetDefault("SENTRY_DSN", "")
viper.SetDefault("METRICS_USERNAME", "") viper.SetDefault("METRICS_USERNAME", "")
viper.SetDefault("METRICS_PASSWORD", "") viper.SetDefault("METRICS_PASSWORD", "")
@@ -80,17 +111,39 @@ 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{ s := &Config{
BindAddress: viper.GetString("BIND_ADDRESS"), BindAddress: viper.GetString("BIND_ADDRESS"),
DataDir: viper.GetString("DATA_DIR"), CORSAllowedOrigins: splitList(viper.GetString("CORS_ALLOWED_ORIGINS")),
Debug: viper.GetBool("DEBUG"), DataDir: viper.GetString("DATA_DIR"),
MetricsPassword: viper.GetString("METRICS_PASSWORD"), DataDirMaxBytes: viper.GetInt64("DATA_DIR_MAX_BYTES"),
MetricsUsername: viper.GetString("METRICS_USERNAME"), Debug: debug,
Port: viper.GetInt("PORT"), MetricsPassword: viper.GetString("METRICS_PASSWORD"),
SentryDSN: viper.GetString("SENTRY_DSN"), MetricsUsername: viper.GetString("METRICS_USERNAME"),
TrustedProxies: splitList(viper.GetString("TRUSTED_PROXIES")), Port: port,
log: log, ReportsPerMinute: viper.GetInt("REPORTS_PER_MINUTE"),
params: &params, SentryDSN: viper.GetString("SENTRY_DSN"),
TrustedProxies: splitList(viper.GetString("TRUSTED_PROXIES")),
log: log,
params: &params,
}
err = s.check()
if err != nil {
return nil, err
} }
if s.Debug { if s.Debug {
@@ -101,6 +154,51 @@ func New(
return s, nil 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.
// go-chi/cors reads a "*" anywhere in an entry as a wildcard, so no
// entry may contain one.
func checkOrigins(origins []string) error {
for _, origin := range origins {
u, err := url.Parse(origin)
if err != nil || u.Scheme == "" || u.Host == "" ||
strings.Contains(origin, "*") ||
origin != u.Scheme+"://"+u.Host {
return fmt.Errorf("CORS_ALLOWED_ORIGINS %q: %w",
origin, errNotOrigin)
}
}
return nil
}
// splitList turns a comma-separated setting into a trimmed // splitList turns a comma-separated setting into a trimmed
// slice, dropping empty entries. // slice, dropping empty entries.
func splitList(raw string) []string { func splitList(raw string) []string {
+121
View File
@@ -0,0 +1,121 @@
package config_test
import (
"strings"
"testing"
"sneak.berlin/go/netwatch/internal/config"
"sneak.berlin/go/netwatch/internal/globals"
"sneak.berlin/go/netwatch/internal/logger"
"go.uber.org/fx"
)
// requireConfigError builds the config as main does and fails the
// test unless that fails with an error naming setting. It uses
// fx.New, because fxtest.New fails the test itself on an error.
func requireConfigError(t *testing.T, setting string) {
t.Helper()
app := fx.New(
fx.NopLogger,
fx.Provide(globals.New, logger.New, config.New),
fx.Invoke(func(*config.Config) {}),
)
err := app.Err()
if err == nil || !strings.Contains(err.Error(), setting) {
t.Fatalf("config error = %v, want one naming %s", err, setting)
}
}
// 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.
func TestReportsPerMinuteMustBePositive(t *testing.T) {
t.Setenv("REPORTS_PER_MINUTE", "0")
requireConfigError(t, "REPORTS_PER_MINUTE")
}
// TestDataDirMaxBytesMustBeANumber: viper reads a value that is not
// a number, such as "1GB", as 0, which would refuse every report.
func TestDataDirMaxBytesMustBeANumber(t *testing.T) {
t.Setenv("DATA_DIR_MAX_BYTES", "1GB")
requireConfigError(t, "DATA_DIR_MAX_BYTES")
}
// TestCORSAllowedOriginsMustBeOrigins: "*" would let every origin in,
// and an entry that is not a plain origin would match no page.
func TestCORSAllowedOriginsMustBeOrigins(t *testing.T) {
for _, entry := range []string{
"*",
"https://*.netwatch.example",
"netwatch.example",
"https://netwatch.example/",
} {
t.Run(entry, func(t *testing.T) {
t.Setenv("CORS_ALLOWED_ORIGINS", entry)
requireConfigError(t, "CORS_ALLOWED_ORIGINS")
})
}
}
+18 -2
View File
@@ -4,6 +4,8 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"net/http" "net/http"
"sneak.berlin/go/netwatch/internal/reportbuf"
) )
// maxLoggedFieldBytes bounds untrusted text (string fields, // maxLoggedFieldBytes bounds untrusted text (string fields,
@@ -55,10 +57,9 @@ func (s *Handlers) HandleReport() http.HandlerFunc {
err = s.buf.Append(rpt) err = s.buf.Append(rpt)
if err != nil { if err != nil {
s.log.Error("failed to buffer report", "error", err)
s.respondJSON(w, r, s.respondJSON(w, r,
&response{Status: "error"}, &response{Status: "error"},
http.StatusInternalServerError, s.appendErrorStatus(err),
) )
return return
@@ -88,6 +89,21 @@ func (s *Handlers) decodeErrorStatus(err error) int {
return http.StatusBadRequest return http.StatusBadRequest
} }
// 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.
func (s *Handlers) appendErrorStatus(err error) int {
if errors.Is(err, reportbuf.ErrFull) {
s.log.Warn("report refused: report files at their size cap")
return http.StatusInsufficientStorage
}
s.log.Error("failed to buffer report", "error", err)
return http.StatusInternalServerError
}
// logReportReceived logs an accepted report. Untrusted fields are // logReportReceived logs an accepted report. Untrusted fields are
// bounded (client_id, timestamp) or reduced to a length // bounded (client_id, timestamp) or reduced to a length
// (geo_bytes) so the raw attacker-controlled body never reaches // (geo_bytes) so the raw attacker-controlled body never reaches
+27
View File
@@ -13,6 +13,7 @@ import (
"sneak.berlin/go/netwatch/internal/handlers" "sneak.berlin/go/netwatch/internal/handlers"
"sneak.berlin/go/netwatch/internal/middleware" "sneak.berlin/go/netwatch/internal/middleware"
"sneak.berlin/go/netwatch/internal/reportbuf"
) )
var errStorageFailed = errors.New("storage failed") var errStorageFailed = errors.New("storage failed")
@@ -66,6 +67,32 @@ 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.
func TestHandleReportFullIs507(t *testing.T) {
t.Parallel()
h := newTestHandlers(stubAppender{err: reportbuf.ErrFull}, io.Discard)
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodPost, "/api/v1/reports",
strings.NewReader(`{"clientId":"c1","hosts":[]}`),
)
h.HandleReport().ServeHTTP(rec, req)
if rec.Code != http.StatusInsufficientStorage {
t.Fatalf("status = %d, want %d",
rec.Code, http.StatusInsufficientStorage)
}
if got := rec.Body.String(); got != "{\"status\":\"error\"}\n" {
t.Errorf("body = %q, want %q", got, "{\"status\":\"error\"}\n")
}
}
func TestHandleReportMalformedJSONIs400(t *testing.T) { func TestHandleReportMalformedJSONIs400(t *testing.T) {
t.Parallel() t.Parallel()
+7 -4
View File
@@ -15,6 +15,13 @@ func NewWithLogger(log *slog.Logger) *Middleware {
return &Middleware{log: log} return &Middleware{log: log}
} }
// NewWithTrustedProxies builds a Middleware that honours forwarded
// headers from the given networks, for tests of the client address
// paths without the fx graph.
func NewWithTrustedProxies(trusted []netip.Prefix) *Middleware {
return &Middleware{trustedProxies: trusted}
}
func ClientIP( func ClientIP(
remoteAddr string, remoteAddr string,
header http.Header, header http.Header,
@@ -22,7 +29,3 @@ func ClientIP(
) string { ) string {
return clientIP(remoteAddr, header, trusted) return clientIP(remoteAddr, header, trusted)
} }
func ParseTrustedProxies(cidrs []string) ([]netip.Prefix, error) {
return parseTrustedProxies(cidrs)
}
+43 -18
View File
@@ -20,6 +20,7 @@ import (
"github.com/go-chi/chi/v5/middleware" "github.com/go-chi/chi/v5/middleware"
"github.com/go-chi/cors" "github.com/go-chi/cors"
"github.com/go-chi/httprate"
"go.uber.org/fx" "go.uber.org/fx"
) )
@@ -63,7 +64,7 @@ func New(
_ fx.Lifecycle, _ fx.Lifecycle,
params Params, params Params,
) (*Middleware, error) { ) (*Middleware, error) {
trusted, err := parseTrustedProxies(params.Config.TrustedProxies) trusted, err := ParseTrustedProxies(params.Config.TrustedProxies)
if err != nil { if err != nil {
return nil, err return nil, err
} }
@@ -76,16 +77,18 @@ func New(
return s, nil return s, nil
} }
// parseTrustedProxies converts CIDR strings into prefixes, // ParseTrustedProxies converts the TRUSTED_PROXIES entries into
// failing fast on any malformed entry. // prefixes, failing fast on any malformed entry. Each entry must be
func parseTrustedProxies(cidrs []string) ([]netip.Prefix, error) { // a CIDR; a lone address is refused. "netwatch-server check-cidr"
// runs it too.
func ParseTrustedProxies(cidrs []string) ([]netip.Prefix, error) {
prefixes := make([]netip.Prefix, 0, len(cidrs)) prefixes := make([]netip.Prefix, 0, len(cidrs))
for _, cidr := range cidrs { for _, cidr := range cidrs {
prefix, err := netip.ParsePrefix(cidr) prefix, err := netip.ParsePrefix(cidr)
if err != nil { if err != nil {
return nil, fmt.Errorf( return nil, fmt.Errorf(
"trusted proxy %q: %w", cidr, err, "TRUSTED_PROXIES %q: %w", cidr, err,
) )
} }
@@ -320,21 +323,43 @@ func (s *Middleware) Recoverer() func(http.Handler) http.Handler {
} }
} }
// CORS returns middleware that adds permissive CORS headers. // CORS returns middleware that lets pages served from the given
func (s *Middleware) CORS() func(http.Handler) http.Handler { // origins call the API. With no origins it adds no CORS headers at
// all, so only same-origin pages can use the API. That case must not
// reach cors.Handler, which treats an empty origin list as "allow
// every origin".
func (s *Middleware) CORS(
origins []string,
) func(http.Handler) http.Handler {
if len(origins) == 0 {
return func(next http.Handler) http.Handler { return next }
}
return cors.Handler(cors.Options{ return cors.Handler(cors.Options{
AllowedOrigins: []string{"*"}, AllowedOrigins: origins,
AllowedMethods: []string{ AllowedMethods: []string{http.MethodGet, http.MethodPost},
"GET", "POST", "PUT", "DELETE", "OPTIONS", AllowedHeaders: []string{"Content-Type"},
},
AllowedHeaders: []string{
"Accept",
"Authorization",
"Content-Type",
"X-CSRF-Token",
},
ExposedHeaders: []string{"Link"},
AllowCredentials: false, AllowCredentials: false,
MaxAge: corsMaxAgeSec, MaxAge: corsMaxAgeSec,
}) })
} }
// RateLimit returns middleware that allows each client address
// perMinute requests a minute and answers the rest with 429, the
// Retry-After header httprate sets, and the usual error body. The
// address is the one clientIP resolves, so clients behind the reverse
// proxy are limited one by one, not together as the proxy.
func (s *Middleware) RateLimit(
perMinute int,
) func(http.Handler) http.Handler {
return httprate.LimitBy(perMinute, time.Minute,
func(r *http.Request) (string, error) {
return clientIP(r.RemoteAddr, r.Header, s.trustedProxies), nil
},
httprate.WithLimitHandler(
func(w http.ResponseWriter, _ *http.Request) {
writeJSONError(w, http.StatusTooManyRequests)
},
),
)
}
+223 -3
View File
@@ -10,6 +10,8 @@ import (
"net/netip" "net/netip"
"strings" "strings"
"testing" "testing"
"testing/synctest"
"time"
"sneak.berlin/go/netwatch/internal/middleware" "sneak.berlin/go/netwatch/internal/middleware"
) )
@@ -34,15 +36,32 @@ func mustPrefixes(t *testing.T, cidrs ...string) []netip.Prefix {
return prefixes return prefixes
} }
// TestParseTrustedProxiesRejectsMalformed includes entries nginx would
// read as another address or look up as a hostname, in the CIDR form
// bin/entrypoint.sh gives "netwatch-server check-cidr".
func TestParseTrustedProxiesRejectsMalformed(t *testing.T) { func TestParseTrustedProxiesRejectsMalformed(t *testing.T) {
t.Parallel() t.Parallel()
_, err := middleware.ParseTrustedProxies([]string{"not-a-cidr"}) for _, cidr := range []string{
if err == nil { "not-a-cidr", "10.0.0.1", "1.2.3/32", "172.30/32", "10/32",
t.Fatal("expected error for malformed CIDR, got nil") "cafe/32", "999.1.1.1/32", "10.0.0.0/33", "::1/129",
"fe80::1%eth0/128",
} {
_, err := middleware.ParseTrustedProxies([]string{cidr})
if err == nil || !strings.Contains(err.Error(), "TRUSTED_PROXIES") {
t.Errorf("%q: error = %v, want one naming TRUSTED_PROXIES",
cidr, err)
}
} }
} }
func TestParseTrustedProxiesAcceptsCIDRs(t *testing.T) {
t.Parallel()
mustPrefixes(t, "172.17.0.1/32", "10.0.0.0/8", "2001:db8::1/128",
"2001:db8::/32", "::ffff:192.0.2.1/128")
}
type clientIPCase struct { type clientIPCase struct {
name string name string
remoteAddr string remoteAddr string
@@ -300,3 +319,204 @@ func TestRecovererRepanicsOnAbortHandler(t *testing.T) {
t.Errorf("abort was logged: %q", logbuf.String()) t.Errorf("abort was logged: %q", logbuf.String())
} }
} }
// okHandler stands in for the route a middleware guards.
func okHandler() http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
})
}
// TestRateLimitRefusesPastAllowanceThenResets checks one client
// address: it may use its whole allowance at once, the next request
// is refused with 429, and later it may send again.
func TestRateLimitRefusesPastAllowanceThenResets(t *testing.T) {
t.Parallel()
// synctest runs this on a fake clock: time.Sleep returns at once,
// with the clock moved on.
synctest.Test(t, func(t *testing.T) {
const perMinute = 2
handler := (&middleware.Middleware{}).RateLimit(perMinute)(okHandler())
post := func() *httptest.ResponseRecorder {
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodPost, "/api/v1/reports", http.NoBody)
handler.ServeHTTP(rec, req)
return rec
}
for i := range perMinute {
if code := post().Code; code != http.StatusOK {
t.Fatalf("request %d: status = %d, want %d",
i+1, code, http.StatusOK)
}
}
rec := post()
if rec.Code != http.StatusTooManyRequests {
t.Fatalf("request past the allowance: status = %d, want %d",
rec.Code, http.StatusTooManyRequests)
}
if got := rec.Body.String(); got != "{\"status\":\"error\"}\n" {
t.Errorf("body = %q, want %q", got, "{\"status\":\"error\"}\n")
}
if got := rec.Header().Get("Retry-After"); got != "60" {
t.Fatalf("Retry-After = %q, want %q", got, "60")
}
// httprate also counts the previous minute's requests, fading
// them out over the current one, so two minutes on the whole
// allowance is back.
time.Sleep(2 * time.Minute)
for i := range perMinute {
if code := post().Code; code != http.StatusOK {
t.Fatalf("two minutes later, request %d: status = %d, want %d",
i+1, code, http.StatusOK)
}
}
})
}
// postForwarded sends handler a report from peer that names client in
// X-Forwarded-For, and returns the status.
func postForwarded(
t *testing.T,
handler http.Handler,
peer, client string,
) int {
t.Helper()
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodPost, "/api/v1/reports", http.NoBody)
req.RemoteAddr = peer
req.Header.Set("X-Forwarded-For", client)
handler.ServeHTTP(rec, req)
return rec.Code
}
// TestRateLimitIsPerForwardedClient checks that clients behind a
// trusted proxy each get their own allowance: the limit is keyed on
// the client address clientIP resolves, not on the proxy's.
func TestRateLimitIsPerForwardedClient(t *testing.T) {
t.Parallel()
const otherClient = "203.0.113.8"
mw := middleware.NewWithTrustedProxies(mustPrefixes(t, "127.0.0.1/32"))
handler := mw.RateLimit(1)(okHandler())
code := postForwarded(t, handler, loopbackPeer, forwardedIP)
if code != http.StatusOK {
t.Fatalf("first request: status = %d, want %d", code, http.StatusOK)
}
code = postForwarded(t, handler, loopbackPeer, forwardedIP)
if code != http.StatusTooManyRequests {
t.Fatalf("same client again: status = %d, want %d",
code, http.StatusTooManyRequests)
}
code = postForwarded(t, handler, loopbackPeer, otherClient)
if code != http.StatusOK {
t.Fatalf("other client behind the same proxy: status = %d, want %d",
code, http.StatusOK)
}
}
// TestRateLimitIgnoresForwardedForFromUntrustedPeer checks that a
// peer that is not a trusted proxy cannot get a fresh allowance by
// naming a different client in X-Forwarded-For on each request.
func TestRateLimitIgnoresForwardedForFromUntrustedPeer(t *testing.T) {
t.Parallel()
const untrustedPeer = "198.51.100.4:5000"
mw := middleware.NewWithTrustedProxies(mustPrefixes(t, "127.0.0.1/32"))
handler := mw.RateLimit(1)(okHandler())
code := postForwarded(t, handler, untrustedPeer, "203.0.113.8")
if code != http.StatusOK {
t.Fatalf("first request: status = %d, want %d", code, http.StatusOK)
}
code = postForwarded(t, handler, untrustedPeer, "203.0.113.9")
if code != http.StatusTooManyRequests {
t.Fatalf("same peer naming another client: status = %d, want %d",
code, http.StatusTooManyRequests)
}
}
// preflight sends cors the preflight request a browser makes before
// it POSTs JSON from origin.
func preflight(
t *testing.T,
cors func(http.Handler) http.Handler,
origin string,
) *httptest.ResponseRecorder {
t.Helper()
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodOptions, "/api/v1/reports", http.NoBody)
req.Header.Set("Origin", origin)
req.Header.Set("Access-Control-Request-Method", http.MethodPost)
req.Header.Set("Access-Control-Request-Headers", "content-type")
cors(okHandler()).ServeHTTP(rec, req)
return rec
}
// TestCORSWithoutOriginsAddsNoHeaders checks the default: with no
// origins configured, no origin is given any CORS header.
func TestCORSWithoutOriginsAddsNoHeaders(t *testing.T) {
t.Parallel()
rec := preflight(t,
(&middleware.Middleware{}).CORS(nil), "https://elsewhere.example")
for name := range rec.Header() {
if strings.HasPrefix(name, "Access-Control-") {
t.Errorf("CORS header %s set with no origins configured", name)
}
}
}
func TestCORSAllowsOnlyListedOrigins(t *testing.T) {
t.Parallel()
const listed = "https://netwatch.example"
cors := (&middleware.Middleware{}).CORS([]string{listed})
cases := []struct {
name string
origin string
want string
}{
{name: "listed origin allowed", origin: listed, want: listed},
{name: "other origin refused", origin: "https://elsewhere.example"},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
rec := preflight(t, cors, tc.origin)
got := rec.Header().Get("Access-Control-Allow-Origin")
if got != tc.want {
t.Errorf("Access-Control-Allow-Origin = %q, want %q",
got, tc.want)
}
})
}
}
+15
View File
@@ -0,0 +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 }
}
+94 -8
View File
@@ -6,12 +6,15 @@ import (
"bytes" "bytes"
"context" "context"
"encoding/json" "encoding/json"
"errors"
"fmt" "fmt"
"io/fs" "io/fs"
"log/slog" "log/slog"
"os" "os"
"path/filepath" "path/filepath"
"strings"
"sync" "sync"
"sync/atomic"
"time" "time"
"sneak.berlin/go/netwatch/internal/config" "sneak.berlin/go/netwatch/internal/config"
@@ -27,8 +30,17 @@ const (
defaultDataDir = "./data/reports" defaultDataDir = "./data/reports"
dirPerms fs.FileMode = 0o750 dirPerms fs.FileMode = 0o750
filePerms fs.FileMode = 0o640 filePerms fs.FileMode = 0o640
// Report files are named filePrefix + timestamp + "-" + number +
// fileSuffix; see writeFile.
filePrefix = "reports-"
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")
// Params defines the dependencies for Buffer. // Params defines the dependencies for Buffer.
type Params struct { type Params struct {
fx.In fx.In
@@ -44,8 +56,19 @@ type Buffer struct {
dataDir string dataDir string
done chan struct{} done chan struct{}
log *slog.Logger log *slog.Logger
maxBytes int64
mu sync.Mutex 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 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.
usedBytes int64
} }
// New creates a Buffer and registers lifecycle hooks to // New creates a Buffer and registers lifecycle hooks to
@@ -60,9 +83,11 @@ func New(
} }
b := &Buffer{ b := &Buffer{
dataDir: dir, dataDir: dir,
done: make(chan struct{}), done: make(chan struct{}),
log: params.Logger.Get(), log: params.Logger.Get(),
maxBytes: params.Config.DataDirMaxBytes,
now: time.Now,
} }
lc.Append(fx.Hook{ lc.Append(fx.Hook{
@@ -72,6 +97,12 @@ 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.
b.usedBytes, err = reportFilesSize(b.dataDir)
if err != nil {
return err
}
go b.flushLoop() go b.flushLoop()
return nil return nil
@@ -97,15 +128,27 @@ 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. If the buffer reaches the size threshold, it is // the buffer. It stores nothing and returns ErrFull if the line
// drained and written to disk asynchronously. // would take usedBytes past maxBytes. If the buffer reaches the
// size threshold, it is drained and written to disk
// 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)
if err != nil { if err != nil {
return fmt.Errorf("marshal report: %w", err) return fmt.Errorf("marshal report: %w", err)
} }
lineBytes := int64(len(line)) + 1 // with its newline
b.mu.Lock() b.mu.Lock()
if b.usedBytes+lineBytes > b.maxBytes {
b.mu.Unlock()
return ErrFull
}
b.usedBytes += lineBytes
b.buf.Write(line) b.buf.Write(line)
b.buf.WriteByte('\n') b.buf.WriteByte('\n')
@@ -177,12 +220,14 @@ func (b *Buffer) drainBuf() []byte {
// writeFile creates a timestamped zstd-compressed JSONL file // writeFile creates a timestamped zstd-compressed JSONL file
// in the data directory. // in the data directory.
func (b *Buffer) writeFile(data []byte) error { func (b *Buffer) writeFile(data []byte) error {
ts := time.Now().UTC().Format("2006-01-02T15-04-05.000Z") // The timestamp comes first, so the names sort by time; the number
name := fmt.Sprintf("reports-%s.jsonl.zst", ts) // 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 := filepath.Join(b.dataDir, name)
// path is built from the operator-supplied dataDir plus a // 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 f, err := os.OpenFile( //nolint:gosec // see comment above
path, path,
os.O_WRONLY|os.O_CREATE|os.O_EXCL, os.O_WRONLY|os.O_CREATE|os.O_EXCL,
@@ -214,10 +259,51 @@ func (b *Buffer) writeFile(data []byte) error {
return fmt.Errorf("close zstd encoder: %w", err) return fmt.Errorf("close zstd encoder: %w", err)
} }
info, err := f.Stat()
if err != nil {
return fmt.Errorf("stat report file: %w", err)
}
err = f.Close() err = f.Close()
if err != nil { if err != nil {
return fmt.Errorf("close report file: %w", err) 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 nil
} }
// reportFilesSize returns the total size of the report files in
// dir.
func reportFilesSize(dir string) (int64, error) {
entries, err := os.ReadDir(dir)
if err != nil {
return 0, fmt.Errorf("read data dir: %w", err)
}
var total int64
for _, entry := range entries {
name := entry.Name()
if !strings.HasPrefix(name, filePrefix) ||
!strings.HasSuffix(name, fileSuffix) {
continue
}
info, err := entry.Info()
if err != nil {
return 0, fmt.Errorf("stat report file: %w", err)
}
total += info.Size()
}
return total, nil
}
@@ -1,17 +1,26 @@
package reportbuf_test package reportbuf_test
import ( import (
"encoding/json"
"errors" "errors"
"fmt"
"io/fs" "io/fs"
"os" "os"
"path/filepath"
"slices"
"strconv"
"strings" "strings"
"sync"
"sync/atomic"
"testing" "testing"
"time"
"sneak.berlin/go/netwatch/internal/config" "sneak.berlin/go/netwatch/internal/config"
"sneak.berlin/go/netwatch/internal/globals" "sneak.berlin/go/netwatch/internal/globals"
"sneak.berlin/go/netwatch/internal/logger" "sneak.berlin/go/netwatch/internal/logger"
"sneak.berlin/go/netwatch/internal/reportbuf" "sneak.berlin/go/netwatch/internal/reportbuf"
"github.com/klauspost/compress/zstd"
"go.uber.org/fx" "go.uber.org/fx"
"go.uber.org/fx/fxtest" "go.uber.org/fx/fxtest"
) )
@@ -94,6 +103,319 @@ func TestFailedFinalFlushFailsStop(t *testing.T) {
} }
} }
// startBuffer starts a Buffer through fx, as main does, with the
// DATA_DIR and DATA_DIR_MAX_BYTES the calling test has set.
func startBuffer(t *testing.T) *reportbuf.Buffer {
t.Helper()
var buf *reportbuf.Buffer
app := fxtest.New(t,
fx.Provide(
globals.New,
logger.New,
config.New,
reportbuf.New,
),
fx.Populate(&buf),
)
app.RequireStart()
t.Cleanup(app.RequireStop)
return buf
}
// lineBytes is what one report takes in the buffer: its JSON and a
// newline.
func lineBytes(t *testing.T, report any) int {
t.Helper()
line, err := json.Marshal(report)
if err != nil {
t.Fatalf("marshal report: %v", err)
}
return len(line) + 1
}
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"}
dir := t.TempDir()
writeBytes(t, filepath.Join(dir, "reports-2026-01-01T00-00-00.000Z.jsonl.zst"),
earlierBytes)
writeBytes(t, filepath.Join(dir, "notes.txt"), 10*earlierBytes)
t.Setenv("DATA_DIR", dir)
t.Setenv("DATA_DIR_MAX_BYTES",
strconv.Itoa(earlierBytes+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)
}
}
// TestWrittenReportsCountAtFileSize checks that once reports are
// written, they count as their compressed file, not their
// uncompressed size, which frees room under the cap.
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)
t.Setenv("DATA_DIR", t.TempDir())
// 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))
buf := startBuffer(t)
err := buf.Append(report)
if err != nil {
t.Fatalf("first report: %v", err)
}
err = buf.Flush()
if err != nil {
t.Fatalf("flush: %v", err)
}
err = buf.Append(report)
if err != nil {
t.Fatalf("second report, after the first was written: %v", err)
}
}
// 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
report := map[string]string{"id": "written"}
size := int64(lineBytes(t, report))
dir := t.TempDir()
t.Setenv("DATA_DIR", dir)
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(maxBytes))
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)
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", used, err)
}
err = buf.Flush()
if err != nil {
t.Fatalf("flush: %v", err)
}
}
t.Fatal("the report files never filled the cap")
}
// TestConcurrentAppendsStopAtCap appends from many goroutines at once
// with room for exactly roomFor reports: exactly that many must be
// taken, which holds only if Append checks and counts each report
// under one lock.
func TestConcurrentAppendsStopAtCap(t *testing.T) {
const (
roomFor = 5
senders = 50
)
// Large, so each Append takes long enough for the senders to
// overlap while the cap is reached.
report := map[string]string{"id": strings.Repeat("a", 1_000_000)}
t.Setenv("DATA_DIR", t.TempDir())
t.Setenv("DATA_DIR_MAX_BYTES",
strconv.Itoa(roomFor*lineBytes(t, report)))
buf := startBuffer(t)
var (
taken atomic.Int64
wg sync.WaitGroup
)
start := make(chan struct{})
for range senders {
wg.Go(func() {
<-start
err := buf.Append(report)
if err == nil {
taken.Add(1)
} else if !errors.Is(err, reportbuf.ErrFull) {
t.Errorf("append: %v", err)
}
})
}
close(start)
wg.Wait()
if got := taken.Load(); got != roomFor {
t.Fatalf("%d reports taken, want %d", got, roomFor)
}
}
// 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()
paths, err := filepath.Glob(filepath.Join(dir, "reports-*.jsonl.zst"))
if err != nil {
t.Fatalf("list report files: %v", err)
}
var total int64
for _, path := range paths {
info, statErr := os.Stat(path)
if statErr != nil {
t.Fatalf("stat %s: %v", path, statErr)
}
total += info.Size()
}
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()
err := os.WriteFile(path, make([]byte, n), 0o600)
if err != nil {
t.Fatalf("write %s: %v", path, err)
}
}
func hasReportFile(t *testing.T, dir string) bool { func hasReportFile(t *testing.T, dir string) bool {
t.Helper() t.Helper()
+3 -2
View File
@@ -25,7 +25,7 @@ func (s *Server) SetupRoutes() {
s.router.Use(middleware.RequestID) s.router.Use(middleware.RequestID)
s.router.Use(s.mw.Logging()) s.router.Use(s.mw.Logging())
s.router.Use(s.mw.SecurityHeaders()) s.router.Use(s.mw.SecurityHeaders())
s.router.Use(s.mw.CORS()) s.router.Use(s.mw.CORS(s.params.Config.CORSAllowedOrigins))
s.router.Use(s.mw.MaxBodyBytes(maxRequestBodyBytes)) s.router.Use(s.mw.MaxBodyBytes(maxRequestBodyBytes))
s.router.Use(middleware.Timeout(requestTimeout)) s.router.Use(middleware.Timeout(requestTimeout))
@@ -35,6 +35,7 @@ func (s *Server) SetupRoutes() {
) )
s.router.Route("/api/v1", func(r chi.Router) { s.router.Route("/api/v1", func(r chi.Router) {
r.Post("/reports", s.h.HandleReport()) r.With(s.mw.RateLimit(s.params.Config.ReportsPerMinute)).
Post("/reports", s.h.HandleReport())
}) })
} }
+58
View File
@@ -49,6 +49,64 @@ func newServer(t *testing.T) *server.Server {
return srv return srv
} }
// TestReportsAreRateLimited checks that POST /api/v1/reports is
// behind the per-address rate limit, set here to two a minute.
func TestReportsAreRateLimited(t *testing.T) {
t.Setenv("REPORTS_PER_MINUTE", "2")
srv := newServer(t)
srv.SetupRoutes()
post := func() int {
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodPost, "/api/v1/reports",
strings.NewReader(`{"clientId":"c1","hosts":[]}`),
)
srv.ServeHTTP(rec, req)
return rec.Code
}
for i := range 2 {
if code := post(); code != http.StatusOK {
t.Fatalf("report %d: status = %d, want %d",
i+1, code, http.StatusOK)
}
}
if code := post(); code != http.StatusTooManyRequests {
t.Fatalf("third report in a minute: status = %d, want %d",
code, http.StatusTooManyRequests)
}
}
// TestCORSAllowedOriginsReachTheRouter checks that an origin listed in
// CORS_ALLOWED_ORIGINS is allowed by the router, not only when handed
// to the CORS middleware directly.
func TestCORSAllowedOriginsReachTheRouter(t *testing.T) {
const origin = "https://netwatch.example:8443"
t.Setenv("CORS_ALLOWED_ORIGINS", origin)
srv := newServer(t)
srv.SetupRoutes()
// The preflight a browser sends before it POSTs JSON from origin.
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodOptions, "/api/v1/reports", http.NoBody)
req.Header.Set("Origin", origin)
req.Header.Set("Access-Control-Request-Method", http.MethodPost)
req.Header.Set("Access-Control-Request-Headers", "content-type")
srv.ServeHTTP(rec, req)
got := rec.Header().Get("Access-Control-Allow-Origin")
if got != origin {
t.Fatalf("Access-Control-Allow-Origin = %q, want %q", got, origin)
}
}
// TestHealthCheckRejectsOversizeBody sends the health check, which // TestHealthCheckRejectsOversizeBody sends the health check, which
// never reads its body, a body one byte over the limit. Only the // never reads its body, a body one byte over the limit. Only the
// router-wide body limit can reject it. // router-wide body limit can reject it.
+66 -6
View File
@@ -8,23 +8,83 @@
# No set -e: kill and wait return non-zero here in normal operation. # No set -e: kill and wait return non-zero here in normal operation.
set -u 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, 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]*)
echo "entrypoint: PORT must be a port number, not '$PORT'" >&2
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
# TRUSTED_PROXIES names the reverse proxies in front of the container,
# as IP addresses or CIDRs separated by commas. nginx takes the client
# address from X-Forwarded-For only on a request from one of them, so
# unset or empty, it trusts no one. nginx.conf includes the file written
# here, one set_real_ip_from line per entry.
#
# nginx looks up an entry it cannot read as an address as a hostname,
# and trusts what it finds (1.2.3 is found as 1.2.0.3). So each entry
# is made a CIDR, a lone address getting /128 if it is IPv6 and /32 if
# not, and netwatch-server checks it with the parsing it gives its own
# TRUSTED_PROXIES. Its error, naming the CIDR, is dropped for the one
# below, naming the entry as written. set -f keeps a * in an entry from
# becoming a list of file names.
TRUSTED_PROXIES="${TRUSTED_PROXIES:-}"
set -f
for proxy in $(printf '%s' "$TRUSTED_PROXIES" | tr ',' ' '); do
case "$proxy" in
*/*) cidr="$proxy" ;;
*:*) cidr="$proxy/128" ;;
*) cidr="$proxy/32" ;;
esac
if ! netwatch-server check-cidr "$cidr" 2> /dev/null; then
echo "entrypoint: TRUSTED_PROXIES must be IP addresses or CIDRs" \
"separated by commas; '$proxy' is neither" >&2
exit 1
fi
echo "set_real_ip_from $cidr;"
done > /etc/nginx/trusted-proxies.conf
# A stop signal is only noted here; the loop below acts on it. # A stop signal is only noted here; the loop below acts on it.
stop_requested="" stop_requested=""
trap 'stop_requested=yes' TERM INT trap 'stop_requested=yes' TERM INT
# netwatch-server runs as the netwatch user and listens on loopback # netwatch-server runs as the netwatch user and listens on loopback
# only, on a port other than the public one; nginx.conf proxies to this # only, on a port other than the public one; nginx.conf proxies to this
# address. The netwatch user has no login shell, hence -s /bin/sh. # address. Its only client is nginx, so it takes the client address
# busybox su replaces itself with the command instead of staying on as # nginx passes on from 127.0.0.1 alone, whatever TRUSTED_PROXIES the
# its parent, so $! is the server's own PID. # container has. The netwatch user has no login shell, hence -s
BIND_ADDRESS=127.0.0.1 PORT=8081 \ # /bin/sh. busybox su replaces itself with the command instead of
# staying on as its parent, so $! is the server's own PID.
BIND_ADDRESS=127.0.0.1 PORT=8081 TRUSTED_PROXIES=127.0.0.1/32 \
su -s /bin/sh netwatch -c 'exec netwatch-server' & su -s /bin/sh netwatch -c 'exec netwatch-server' &
backend=$! backend=$!
# nginx starts through the nginx image's own entrypoint, which applies # nginx starts through the nginx image's own entrypoint, which applies
# the image's start-up configuration and then replaces itself with # the image's start-up configuration and then replaces itself with
# nginx. # nginx. Part of that start-up configuration renders nginx.conf into
/docker-entrypoint.sh nginx -g 'daemon off;' & # conf.d with nginx listening on PORT. NGINX_ENVSUBST_FILTER limits
# that rendering to PORT: a variable nginx itself uses, such as $uri,
# would otherwise be replaced by an environment variable of the same
# name.
NGINX_ENVSUBST_FILTER='^PORT$' \
/docker-entrypoint.sh nginx -g 'daemon off;' &
nginx=$! nginx=$!
running() { running() {
+13 -5
View File
@@ -1,14 +1,22 @@
# A template: the nginx image renders it into conf.d at container start,
# filling in PORT and nothing else. bin/entrypoint.sh sets PORT and that
# limit.
server { server {
listen 8080; listen ${PORT};
server_name _; server_name _;
# Keep the nginx version out of the Server header and error pages.
server_tokens off;
root /usr/share/nginx/html; root /usr/share/nginx/html;
index index.html; index index.html;
# Trust RFC1918 reverse proxies for X-Forwarded-For # The client address comes from X-Forwarded-For only on a request
set_real_ip_from 10.0.0.0/8; # from the reverse proxies in TRUSTED_PROXIES: bin/entrypoint.sh
set_real_ip_from 172.16.0.0/12; # writes one set_real_ip_from line for each into this file, and
set_real_ip_from 192.168.0.0/16; # leaves it empty when TRUSTED_PROXIES is unset, so that by default
# the client address is the one each request comes from.
include /etc/nginx/trusted-proxies.conf;
real_ip_header X-Forwarded-For; real_ip_header X-Forwarded-For;
real_ip_recursive on; real_ip_recursive on;
+7 -1
View File
@@ -61,10 +61,16 @@ main() {
# host. # host.
docker network create --internal "$NETWORK" > /dev/null docker network create --internal "$NETWORK" > /dev/null
# nginx.conf is a template: the image renders it over its own
# default.conf, with the same port and limit bin/entrypoint.sh uses.
# The empty file it includes trusts no proxy, as bin/entrypoint.sh
# writes it when TRUSTED_PROXIES is unset.
docker run -d --rm --name "$SERVER" \ docker run -d --rm --name "$SERVER" \
--network "$NETWORK" --network-alias netwatch \ --network "$NETWORK" --network-alias netwatch \
-e PORT=8080 -e NGINX_ENVSUBST_FILTER='^PORT$' \
-v "$ROOT/dist:/usr/share/nginx/html:ro" \ -v "$ROOT/dist:/usr/share/nginx/html:ro" \
-v "$ROOT/nginx.conf:/etc/nginx/conf.d/default.conf:ro" \ -v "$ROOT/nginx.conf:/etc/nginx/templates/default.conf.template:ro" \
-v /dev/null:/etc/nginx/trusted-proxies.conf:ro \
"$SERVER_IMAGE" > /dev/null "$SERVER_IMAGE" > /dev/null
# The image's own entrypoint already exposes CDP on 9222 and passes # The image's own entrypoint already exposes CDP on 9222 and passes