Compare commits
6
Commits
9d0462a2a6
..
next
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8833603eff | ||
|
|
6022cc8b02 | ||
|
|
d2f219ca19 | ||
|
|
ced1956b06 | ||
|
|
ea66caf338 | ||
|
|
bbcc7d921d |
+12
-3
@@ -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
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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
@@ -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
|
||||||
|
|
||||||
|
|||||||
@@ -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
@@ -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
@@ -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=
|
||||||
|
|||||||
@@ -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: ¶ms,
|
SentryDSN: viper.GetString("SENTRY_DSN"),
|
||||||
|
TrustedProxies: splitList(viper.GetString("TRUSTED_PROXIES")),
|
||||||
|
log: log,
|
||||||
|
params: ¶ms,
|
||||||
|
}
|
||||||
|
|
||||||
|
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 {
|
||||||
|
|||||||
@@ -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")
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
|||||||
@@ -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()
|
||||||
|
|
||||||
|
|||||||
@@ -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)
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -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)
|
||||||
|
},
|
||||||
|
),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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 }
|
||||||
|
}
|
||||||
@@ -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()
|
||||||
|
|
||||||
|
|||||||
@@ -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())
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
@@ -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
@@ -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;
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
Reference in New Issue
Block a user