1 Commits
Author SHA1 Message Date
sneak 02a503b9d7 feat(frontend): post collected samples to /api/v1/reports (closes #53)
check / check (push) Failing after 1s
A Reporter beside AppState POSTs a JSON delta report to the same-origin
/api/v1/reports every reportInterval (default 60s). buildReport is an
exported pure function of host state, emitting each host's unreported,
non-paused samples plus a per-browser clientId, geo null, and a UTC
timestamp; a per-host high-water mark advances only on a delivered POST,
so a failed send re-sends next interval within the history window.
Failure is quiet and never blocks probing.

The client id feature-detects crypto.randomUUID and otherwise builds a v4
id from crypto.getRandomValues, so insecure-context loads (plain HTTP to a
non-localhost host) work; reporter setup is isolated so it can never stop
probing. init() runs only when the #app page is present, so a test can
import buildReport.

vite.config.js proxies /api to 127.0.0.1:8080 for yarn dev.

Model: opus-4-8
2026-09-21 16:49:23 +00:00
17 changed files with 95 additions and 838 deletions
+2 -2
View File
@@ -1,8 +1,8 @@
# Lint stage — fast feedback on formatting and lint issues. The
# golangci/golangci-lint image ships Go, gofmt, make and the linter, so
# nothing is installed here.
# golangci/golangci-lint:v2.12.2 (2026-08-10)
FROM golangci/golangci-lint@sha256:5cceeef04e53efe1470638d4b4b4f5ceefd574955ab3941b2d9a68a8c9ad5240 AS lint
# golangci/golangci-lint:v2.7.2 (2026-08-09)
FROM golangci/golangci-lint@sha256:5d6d5c70a61f1356adfd9dd6316ce286799fefc9d743421356ff1b00842368ba AS lint
WORKDIR /src
COPY backend/go.mod backend/go.sum ./
+9 -36
View File
@@ -10,29 +10,18 @@
# Status
pre-1.0. No git tags. `feat/reportbuf-storage` is merged; the backend, the CI
workflow, and the backend repo standard files are all on `main`. Frontend and
backend are both functional. Working toward the 1.0.0 milestone by closing the
remaining repo-compliance issues on the tracker.
pre-1.0. No git tags. Backend work in flight on feat/reportbuf-storage (dirty:
src/main.js). Frontend is functional; backend is new and unmerged.
# Next Step
Confirm the `.gitea/workflows/check.yml` run is green (main always green
policy). The workflow file is already on `main`; what is unverified is that its
latest run passes.
Land feat/reportbuf-storage: finish the in-progress src/main.js change, get make
check green, and merge the branch to main. The branch adds the backend (buffered
zstd-compressed report storage), the CI workflow, and backend repo standard
files, so merging it also closes most compliance gaps.
# Completed Steps
- 2026-09-28: report ingest correctness (issue #23): a storage failure now
returns 500 instead of a false `ok`; oversize bodies return 413 (distinguished
from malformed JSON, which stays 400); a `MaxBodyBytes` middleware caps every
route, not just the report route; the raw attacker-controlled `geo` blob is no
longer logged (only its length), and `client_id`, `timestamp` and decode error
text are length-bounded before logging; a `decodeJSON` handler helper was
added; panic recovery now routes the stack through slog instead of chi's
plain-text stderr; and writing a report file now returns its error, so a
failed final flush on shutdown makes the process exit non-zero instead of
losing the buffered reports silently
- 2026-09-21: shutdown lifecycle correctness. The process now shuts down through
fx instead of `os.Exit`, so every component's `OnStop` runs and buffered
reports are flushed to disk on `SIGTERM` — previously a full flush window of
@@ -45,23 +34,13 @@ latest run passes.
- 2026-09-21: frontend reporting client — a `Reporter` class posts collected
samples to `/api/v1/reports` every `reportInterval` (default 60s) as a
per-host delta, with the report-building step a pure exported function of host
state; only one report POST is in flight at a time and it is abandoned after
half the interval, so a slow backend cannot cause re-sent samples or a mark
moving backwards; the per-browser client id works in insecure (plain-HTTP)
contexts; `vite.config.js` proxies `/api` to the local backend for `yarn dev`
state; the per-browser client id works in insecure (plain-HTTP) contexts;
`vite.config.js` proxies `/api` to the local backend for `yarn dev`
- 2026-09-21: backend HTTP hardening (issue #19): added `ReadHeaderTimeout` and
`IdleTimeout` to the server, a `SecurityHeaders` middleware (HSTS, tight CSP,
frame/sniff/referrer/permissions headers) registered before CORS, and
trusted-proxy client IP resolution honouring `X-Forwarded-For` / `X-Real-IP`
only from a `TRUSTED_PROXIES` allowlist (loopback plus RFC1918 by default)
- 2026-08-10: adopted the org-standard `backend/.golangci.yml` verbatim and
moved the pinned golangci-lint from v2.7.2 to v2.12.2 (the `lint` stage of
`Dockerfile.backend` now pins the `golangci/golangci-lint:v2.12.2` image by
digest); the previous config declared `version: "2"` but used v1 schema keys,
so every threshold in it was inert and its green result was meaningless.
`backend/Makefile`'s `lint` target now asserts the config's sha256 against the
canonical file first, so drift from the org standard fails the build instead
of silently degrading to defaults
- 2026-08-10: every interactive control now meets the 44x44 CSS px minimum tap
target (`.pin-btn`, `#interval-select`, the debug-log label and, on narrow
viewports, `#pause-btn`). The pin button's hit area grows via matching
@@ -91,7 +70,7 @@ latest run passes.
shims, README Entrypoints section
- 2026-02-27: backend with buffered zstd-compressed report storage; CI workflow
and backend repo standard files; backend Dockerfile fixed (Go 1.25,
golangci-lint) and moved to repo root (feat/reportbuf-storage)
golangci-lint) and moved to repo root (feat/reportbuf-storage, unmerged)
- 2026-02-26: host row layout redesigned with CSS grid; overflow and spacing
fixes; nginx config extracted; port hardcoded to 8080
- 2026-02-26: debug log panel, median stats, recovery probe, Docker build fix,
@@ -113,9 +92,3 @@ latest run passes.
(main always green policy)
- Decide what to do with untracked resume.sh: commit it, gitignore it, or delete
it
- Upstream fix needed in `sneak/prompts`: the org-standard `.golangci.yml`
enables `gomodguard`, which golangci-lint v2.12.2 reports as deprecated since
v2.12.0 and replaced by `gomodguard_v2`, so every backend lint run prints a
deprecation warning. The file is standardized and must never be edited in this
repo, so nothing can be done here beyond tracking it — tracked at
<https://git.eeqj.de/sneak/netwatch/issues/41>
+3 -5
View File
@@ -1,9 +1,5 @@
version: "2"
# Config schema uses the golangci-lint v2 layout (settings live under
# linters.settings, not top-level linters-settings) so that the
# thresholds below are actually applied by golangci-lint >= v2.
run:
timeout: 5m
modules-download-mode: readonly
@@ -18,7 +14,8 @@ linters:
- wsl # Deprecated, replaced by wsl_v5
- wrapcheck # Too verbose for internal packages
- varnamelen # Short names like db, id are idiomatic Go
settings:
linters-settings:
lll:
line-length: 88
funlen:
@@ -30,5 +27,6 @@ linters:
threshold: 100
issues:
exclude-use-default: false
max-issues-per-linter: 0
max-same-issues: 0
-19
View File
@@ -10,17 +10,6 @@ GOLDFLAGS += -s -w
GOLDFLAGS += -X main.Version=$(VERSION)
GOLDFLAGS += -X main.Buildarch=$(BUILDARCH)
# macOS ships shasum rather than sha256sum.
SHA256SUM := $(shell command -v sha256sum >/dev/null 2>&1 && echo sha256sum || echo shasum -a 256)
# .golangci.yml is standardized org-wide and must never be edited here
# (REPO_POLICIES.md). Its last silent drift replaced the v2 schema with
# v1 keys, which left every threshold in the file inert while the build
# stayed green. The lint target therefore asserts the file still matches
# the canonical copy byte for byte. The check is a local hash comparison:
# no network, no remote schema, nothing unpinned in the build path.
GOLANGCI_CONFIG_SHA256 := 021cc83f4e6fc7c31b95b34b846723dfcf20b66b7baeea1dc40406e643346bcb
.PHONY: all build test lint fmt fmt-check check docker hooks run clean
all: build
@@ -33,14 +22,6 @@ test:
timeout 30 go test ./...
lint:
@actual=$$($(SHA256SUM) .golangci.yml | cut -d' ' -f1); \
if [ "$$actual" != "$(GOLANGCI_CONFIG_SHA256)" ]; then \
echo ".golangci.yml has drifted from the org standard."; \
echo " expected $(GOLANGCI_CONFIG_SHA256)"; \
echo " actual $$actual"; \
echo "Restore it verbatim from sneak/prompts; do not edit it."; \
exit 1; \
fi
golangci-lint run ./...
fmt:
-13
View File
@@ -1,13 +0,0 @@
package handlers
import "log/slog"
// MaxLoggedFieldBytes exposes the log bound to the external tests.
const MaxLoggedFieldBytes = maxLoggedFieldBytes
// NewForTest builds a Handlers around a report sink and logger,
// bypassing the fx graph so handler behaviour (including the
// storage failure path) is exercisable in unit tests.
func NewForTest(buf reportAppender, log *slog.Logger) *Handlers {
return &Handlers{buf: buf, log: log}
}
+1 -20
View File
@@ -18,13 +18,6 @@ import (
const jsonContentType = "application/json; charset=utf-8"
// reportAppender is the subset of the report buffer the handlers
// depend on. Defining it here keeps the storage failure path
// exercisable with a stub in tests.
type reportAppender interface {
Append(v any) error
}
// Params defines the dependencies for Handlers.
type Params struct {
fx.In
@@ -37,7 +30,7 @@ type Params struct {
// Handlers provides HTTP handler factories for all endpoints.
type Handlers struct {
buf reportAppender
buf *reportbuf.Buffer
hc *healthcheck.Healthcheck
log *slog.Logger
params *Params
@@ -79,15 +72,3 @@ func (s *Handlers) respondJSON(
}
}
}
// decodeJSON decodes the request body into v. The body is
// expected to already be bounded by the body-size middleware, so
// a caller can distinguish an over-limit body from malformed
// JSON by testing the returned error for *http.MaxBytesError.
func (s *Handlers) decodeJSON(
_ http.ResponseWriter,
r *http.Request,
v any,
) error {
return json.NewDecoder(r.Body).Decode(v)
}
+25 -61
View File
@@ -2,14 +2,10 @@ package handlers
import (
"encoding/json"
"errors"
"net/http"
)
// maxLoggedFieldBytes bounds untrusted text (string fields,
// decode error text) before it is logged, so a caller cannot
// inflate log volume with an oversized value.
const maxLoggedFieldBytes = 128
const maxReportBodyBytes = 1 << 20 // 1 MiB
type reportSample struct {
T int64 `json:"t"`
@@ -39,80 +35,48 @@ func (s *Handlers) HandleReport() http.HandlerFunc {
}
return func(w http.ResponseWriter, r *http.Request) {
r.Body = http.MaxBytesReader(
w, r.Body, maxReportBodyBytes,
)
var rpt report
err := s.decodeJSON(w, r, &rpt)
err := json.NewDecoder(r.Body).Decode(&rpt)
if err != nil {
s.respondJSON(w, r,
&response{Status: "error"},
s.decodeErrorStatus(err),
)
return
}
s.logReportReceived(rpt)
err = s.buf.Append(rpt)
if err != nil {
s.log.Error("failed to buffer report", "error", err)
s.respondJSON(w, r,
&response{Status: "error"},
http.StatusInternalServerError,
)
return
}
s.respondJSON(w, r, &response{Status: "ok"}, http.StatusOK)
}
}
// decodeErrorStatus logs a report decode failure and returns the
// status to send: 413 when the body exceeded the size limit,
// otherwise 400 for malformed JSON.
func (s *Handlers) decodeErrorStatus(err error) int {
var tooLarge *http.MaxBytesError
if errors.As(err, &tooLarge) {
s.log.Warn("report body too large", "limit_bytes", tooLarge.Limit)
return http.StatusRequestEntityTooLarge
}
// The decoder's error text can quote request bytes (a whole
// oversized number, for example), so it is bounded too.
s.log.Error("failed to decode report",
"error", boundedForLog(err.Error()),
"error", err,
)
s.respondJSON(w, r,
&response{Status: "error"},
http.StatusBadRequest,
)
return http.StatusBadRequest
return
}
// logReportReceived logs an accepted report. Untrusted fields are
// bounded (client_id, timestamp) or reduced to a length
// (geo_bytes) so the raw attacker-controlled body never reaches
// the log.
func (s *Handlers) logReportReceived(rpt report) {
totalSamples := 0
for _, h := range rpt.Hosts {
totalSamples += len(h.History)
}
s.log.Info("report received",
"client_id", boundedForLog(rpt.ClientID),
"timestamp", boundedForLog(rpt.Timestamp),
"client_id", rpt.ClientID,
"timestamp", rpt.Timestamp,
"host_count", len(rpt.Hosts),
"total_samples", totalSamples,
"geo_bytes", len(rpt.Geo),
"geo", string(rpt.Geo),
)
bufErr := s.buf.Append(rpt)
if bufErr != nil {
s.log.Error("failed to buffer report",
"error", bufErr,
)
}
// boundedForLog truncates an untrusted string to a fixed byte
// bound so an attacker-controlled field cannot dominate the log.
func boundedForLog(s string) string {
if len(s) > maxLoggedFieldBytes {
return s[:maxLoggedFieldBytes]
s.respondJSON(w, r,
&response{Status: "ok"},
http.StatusOK,
)
}
return s
}
-212
View File
@@ -1,212 +0,0 @@
package handlers_test
import (
"bytes"
"encoding/json"
"errors"
"io"
"log/slog"
"net/http"
"net/http/httptest"
"strings"
"testing"
"sneak.berlin/go/netwatch/internal/handlers"
"sneak.berlin/go/netwatch/internal/middleware"
)
var errStorageFailed = errors.New("storage failed")
// stubAppender drives the storage success/failure path without a
// real buffer or disk.
type stubAppender struct {
err error
}
func (s stubAppender) Append(any) error { return s.err }
func newTestHandlers(buf stubAppender, out io.Writer) *handlers.Handlers {
return handlers.NewForTest(buf, slog.New(slog.NewJSONHandler(out, nil)))
}
func decodeStatus(t *testing.T, body []byte) string {
t.Helper()
var resp struct {
Status string `json:"status"`
}
err := json.Unmarshal(body, &resp)
if err != nil {
t.Fatalf("response body not JSON: %v (%q)", err, body)
}
return resp.Status
}
func TestHandleReportStorageFailureIsNon2xx(t *testing.T) {
t.Parallel()
h := newTestHandlers(stubAppender{err: errStorageFailed}, 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 < 500 {
t.Fatalf("storage failure status = %d, want a 5xx", rec.Code)
}
if got := decodeStatus(t, rec.Body.Bytes()); got != "error" {
t.Fatalf("status field = %q, want %q", got, "error")
}
}
func TestHandleReportMalformedJSONIs400(t *testing.T) {
t.Parallel()
h := newTestHandlers(stubAppender{}, io.Discard)
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodPost, "/api/v1/reports",
strings.NewReader(`{not json`),
)
h.HandleReport().ServeHTTP(rec, req)
if rec.Code != http.StatusBadRequest {
t.Fatalf("malformed status = %d, want %d",
rec.Code, http.StatusBadRequest)
}
}
func TestHandleReportOversizeIs413(t *testing.T) {
t.Parallel()
const limit = 32
h := newTestHandlers(stubAppender{}, io.Discard)
handler := (&middleware.Middleware{}).MaxBodyBytes(limit)(
h.HandleReport(),
)
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodPost, "/api/v1/reports",
strings.NewReader(`{"clientId":"`+strings.Repeat("x", 200)+`"}`),
)
// No declared length, so only the middleware's read cap can
// stop this body.
req.ContentLength = -1
handler.ServeHTTP(rec, req)
if rec.Code != http.StatusRequestEntityTooLarge {
t.Fatalf("oversize status = %d, want %d",
rec.Code, http.StatusRequestEntityTooLarge)
}
}
func TestHandleReportDoesNotLogRawGeo(t *testing.T) {
t.Parallel()
const sentinel = "SENSITIVE-GEO-BLOB"
var logbuf bytes.Buffer
h := newTestHandlers(stubAppender{}, &logbuf)
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodPost, "/api/v1/reports",
strings.NewReader(
`{"clientId":"c1","geo":{"raw":"`+sentinel+`"},"hosts":[]}`,
),
)
h.HandleReport().ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusOK)
}
if strings.Contains(logbuf.String(), sentinel) {
t.Fatal("raw geo bytes were written to the log")
}
if !strings.Contains(logbuf.String(), "geo_bytes") {
t.Fatal("expected a bounded geo_bytes field in the log")
}
}
func TestHandleReportLogsClientIDCutToBound(t *testing.T) {
t.Parallel()
long := strings.Repeat("c", 2*handlers.MaxLoggedFieldBytes)
var logbuf bytes.Buffer
h := newTestHandlers(stubAppender{}, &logbuf)
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodPost, "/api/v1/reports",
strings.NewReader(
`{"clientId":"`+long+`","timestamp":"`+long+`","hosts":[]}`,
),
)
h.HandleReport().ServeHTTP(rec, req)
var logged map[string]any
err := json.Unmarshal(logbuf.Bytes(), &logged)
if err != nil {
t.Fatalf("log line not JSON: %v (%q)", err, logbuf.String())
}
want := long[:handlers.MaxLoggedFieldBytes]
if logged["client_id"] != want {
t.Fatalf("logged client_id not cut to %d bytes: %q",
handlers.MaxLoggedFieldBytes, logged["client_id"])
}
if logged["timestamp"] != want {
t.Fatalf("logged timestamp not cut to %d bytes: %q",
handlers.MaxLoggedFieldBytes, logged["timestamp"])
}
}
func TestHandleReportDecodeErrorLogIsBounded(t *testing.T) {
t.Parallel()
// A number too large for its int64 field makes the decoder's
// error text quote the whole number.
huge := strings.Repeat("9", 2*handlers.MaxLoggedFieldBytes)
var logbuf bytes.Buffer
h := newTestHandlers(stubAppender{}, &logbuf)
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodPost, "/api/v1/reports",
strings.NewReader(`{"hosts":[{"history":[{"t":`+huge+`}]}]}`),
)
h.HandleReport().ServeHTTP(rec, req)
if rec.Code != http.StatusBadRequest {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusBadRequest)
}
if strings.Contains(logbuf.String(), huge) {
t.Fatal("the whole oversized number was written to the log")
}
}
@@ -1,7 +1,6 @@
package middleware
import (
"log/slog"
"net/http"
"net/netip"
)
@@ -9,12 +8,6 @@ import (
// Test-only wrappers exposing unexported helpers to the
// external middleware_test package.
// NewWithLogger builds a Middleware around a logger for tests
// that exercise the logging paths without the fx graph.
func NewWithLogger(log *slog.Logger) *Middleware {
return &Middleware{log: log}
}
func ClientIP(
remoteAddr string,
header http.Header,
-84
View File
@@ -3,14 +3,11 @@
package middleware
import (
"errors"
"fmt"
"io"
"log/slog"
"net"
"net/http"
"net/netip"
"runtime/debug"
"strings"
"time"
@@ -25,14 +22,6 @@ import (
const corsMaxAgeSec = 300
// jsonErrorBody is the body written for errors raised inside
// middleware, matching the {"status":"error"} shape the handlers
// return so clients see one error contract across the API.
const (
jsonContentType = "application/json; charset=utf-8"
jsonErrorBody = "{\"status\":\"error\"}\n"
)
// Security header values. The backend is a JSON API with no
// HTML surface, so the CSP forbids every resource type and
// framing outright.
@@ -247,79 +236,6 @@ func (s *Middleware) SecurityHeaders() func(http.Handler) http.Handler {
}
}
// writeJSONError writes the shared JSON error body with the
// given status. Used where middleware must reject a request
// before it reaches a handler.
func writeJSONError(w http.ResponseWriter, status int) {
w.Header().Set("Content-Type", jsonContentType)
w.WriteHeader(status)
_, _ = io.WriteString(w, jsonErrorBody)
}
// MaxBodyBytes returns middleware that caps the request body at
// limit bytes. A declared Content-Length over the limit is
// rejected immediately with 413. Bodies without a declared
// length (or that understate it) are capped as they are read, so
// a handler that reads the body sees a *http.MaxBytesError it can
// map to 413. Mounted again on a route group, it can only lower
// the limit: a cap applied earlier in the chain still holds.
func (s *Middleware) MaxBodyBytes(
limit int64,
) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
if r.ContentLength > limit {
writeJSONError(
w,
http.StatusRequestEntityTooLarge,
)
return
}
r.Body = http.MaxBytesReader(w, r.Body, limit)
next.ServeHTTP(w, r)
},
)
}
}
// Recoverer returns middleware that recovers from a panic in a
// downstream handler, logs the panic and stack trace through
// slog, and responds 500 with no body. http.ErrAbortHandler is
// re-panicked so the server can abort the response as intended.
func (s *Middleware) Recoverer() func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
defer func() {
rec := recover()
if rec == nil {
return
}
err, ok := rec.(error)
if ok && errors.Is(err, http.ErrAbortHandler) {
panic(rec)
}
s.log.ErrorContext(r.Context(),
"panic recovered",
"panic", fmt.Sprintf("%v", rec),
"stack", string(debug.Stack()),
)
w.WriteHeader(http.StatusInternalServerError)
}()
next.ServeHTTP(w, r)
},
)
}
}
// CORS returns middleware that adds permissive CORS headers.
func (s *Middleware) CORS() func(http.Handler) http.Handler {
return cors.Handler(cors.Options{
+12 -175
View File
@@ -1,28 +1,14 @@
package middleware_test
import (
"bytes"
"encoding/json"
"errors"
"log/slog"
"net/http"
"net/http/httptest"
"net/netip"
"strings"
"testing"
"sneak.berlin/go/netwatch/internal/middleware"
)
const (
// loopbackPeer is a remote address inside the trusted-proxy allowlist.
loopbackPeer = "127.0.0.1:5000"
// forwardedIP is the client address presented via X-Forwarded-For.
forwardedIP = "203.0.113.7"
// realIP is the client address presented via X-Real-IP.
realIP = "203.0.113.9"
)
func mustPrefixes(t *testing.T, cidrs ...string) []netip.Prefix {
t.Helper()
@@ -55,32 +41,32 @@ func clientIPCases() []clientIPCase {
return []clientIPCase{
{
name: "trusted proxy uses forwarded-for",
remoteAddr: loopbackPeer,
xff: forwardedIP,
want: forwardedIP,
remoteAddr: "127.0.0.1:5000",
xff: "203.0.113.7",
want: "203.0.113.7",
},
{
name: "trusted proxy uses left-most of chain",
remoteAddr: "10.1.2.3:5000",
xff: forwardedIP + ", 10.1.2.3",
want: forwardedIP,
xff: "203.0.113.7, 10.1.2.3",
want: "203.0.113.7",
},
{
name: "trusted proxy falls back to x-real-ip",
remoteAddr: loopbackPeer,
xRealIP: realIP,
want: realIP,
remoteAddr: "127.0.0.1:5000",
xRealIP: "203.0.113.9",
want: "203.0.113.9",
},
{
name: "untrusted peer ignores forwarded-for",
remoteAddr: "198.51.100.4:5000",
xff: forwardedIP,
xff: "203.0.113.7",
want: "198.51.100.4",
},
{
name: "untrusted peer ignores x-real-ip",
remoteAddr: "198.51.100.4:5000",
xRealIP: realIP,
xRealIP: "203.0.113.9",
want: "198.51.100.4",
},
{
@@ -90,7 +76,7 @@ func clientIPCases() []clientIPCase {
},
{
name: "trusted proxy with garbage header uses peer",
remoteAddr: loopbackPeer,
remoteAddr: "127.0.0.1:5000",
xff: "not-an-ip",
want: "127.0.0.1",
},
@@ -133,7 +119,7 @@ func TestSecurityHeaders(t *testing.T) {
)
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(), http.MethodGet, "/", http.NoBody)
req := httptest.NewRequest(http.MethodGet, "/", http.NoBody)
handler.ServeHTTP(rec, req)
want := map[string]string{
@@ -151,152 +137,3 @@ func TestSecurityHeaders(t *testing.T) {
}
}
}
// TestMaxBodyBytesRejectsOversizeOnNonReadingRoute confirms the
// limit is enforced even for a handler that never reads the body
// (for example the health check), via the Content-Length check.
func TestMaxBodyBytesRejectsOversizeOnNonReadingRoute(t *testing.T) {
t.Parallel()
const limit = 16
called := false
handler := (&middleware.Middleware{}).MaxBodyBytes(limit)(
http.HandlerFunc(func(_ http.ResponseWriter, _ *http.Request) {
called = true
}),
)
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodPost, "/.well-known/healthcheck",
strings.NewReader(strings.Repeat("x", limit+1)),
)
handler.ServeHTTP(rec, req)
if rec.Code != http.StatusRequestEntityTooLarge {
t.Fatalf("status = %d, want %d",
rec.Code, http.StatusRequestEntityTooLarge)
}
if called {
t.Fatal("handler ran despite oversize body")
}
if got := rec.Body.String(); got != "{\"status\":\"error\"}\n" {
t.Errorf("body = %q, want %q", got, "{\"status\":\"error\"}\n")
}
got := rec.Header().Get("Content-Type")
if got != "application/json; charset=utf-8" {
t.Errorf("Content-Type = %q, want a JSON content type", got)
}
}
func TestMaxBodyBytesAllowsWithinLimit(t *testing.T) {
t.Parallel()
const limit = 64
handler := (&middleware.Middleware{}).MaxBodyBytes(limit)(
http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
}),
)
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodPost, "/api/v1/reports",
strings.NewReader(`{"clientId":"c1"}`),
)
handler.ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusOK)
}
}
func TestRecovererReturns500AndLogsThroughSlog(t *testing.T) {
t.Parallel()
var logbuf bytes.Buffer
mw := middleware.NewWithLogger(
slog.New(slog.NewJSONHandler(&logbuf, nil)),
)
handler := mw.Recoverer()(
http.HandlerFunc(func(_ http.ResponseWriter, _ *http.Request) {
panic("boom")
}),
)
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(), http.MethodGet, "/", http.NoBody)
handler.ServeHTTP(rec, req)
if rec.Code != http.StatusInternalServerError {
t.Fatalf("status = %d, want %d",
rec.Code, http.StatusInternalServerError)
}
var record map[string]any
err := json.Unmarshal(logbuf.Bytes(), &record)
if err != nil {
t.Fatalf("panic log is not one JSON record: %v (%q)", err, logbuf.String())
}
if record["msg"] != "panic recovered" || record["level"] != "ERROR" {
t.Errorf("log record = %v, want msg %q at level ERROR",
record, "panic recovered")
}
if record["panic"] != "boom" {
t.Errorf("panic field = %v, want %q", record["panic"], "boom")
}
stack, _ := record["stack"].(string)
if !strings.HasPrefix(stack, "goroutine ") {
t.Errorf("stack field = %q, want a stack trace", stack)
}
}
// TestRecovererRepanicsOnAbortHandler checks that a handler aborting
// with http.ErrAbortHandler is not treated as a crash: Recoverer
// panics again so the server aborts the response, and logs nothing.
func TestRecovererRepanicsOnAbortHandler(t *testing.T) {
t.Parallel()
var logbuf bytes.Buffer
mw := middleware.NewWithLogger(
slog.New(slog.NewJSONHandler(&logbuf, nil)),
)
handler := mw.Recoverer()(
http.HandlerFunc(func(_ http.ResponseWriter, _ *http.Request) {
panic(http.ErrAbortHandler)
}),
)
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(), http.MethodGet, "/", http.NoBody)
var recovered any
func() {
defer func() { recovered = recover() }()
handler.ServeHTTP(rec, req)
}()
err, _ := recovered.(error)
if !errors.Is(err, http.ErrAbortHandler) {
t.Errorf("Recoverer panicked with %v, want http.ErrAbortHandler", recovered)
}
if logbuf.Len() != 0 {
t.Errorf("abort was logged: %q", logbuf.String())
}
}
+23 -41
View File
@@ -80,16 +80,12 @@ func New(
// stopOnce makes OnStop idempotent: a second
// invocation must not close an already-closed channel
// (which would panic) or flush again.
var err error
b.stopOnce.Do(func() {
close(b.done)
err = b.flushLocked()
b.flushLocked()
})
// A failed final flush fails the stop, so the process
// exits non-zero.
return err
return nil
},
})
@@ -113,12 +109,7 @@ func (b *Buffer) Append(v any) error {
data := b.drainBuf()
b.mu.Unlock()
go func() {
writeErr := b.writeFile(data)
if writeErr != nil {
b.log.Error("flush reports failed", "error", writeErr)
}
}()
go b.writeFile(data)
return nil
}
@@ -137,10 +128,7 @@ func (b *Buffer) flushLoop() {
for {
select {
case <-ticker.C:
err := b.flushLocked()
if err != nil {
b.log.Error("flush reports failed", "error", err)
}
b.flushLocked()
case <-b.done:
return
}
@@ -149,19 +137,19 @@ func (b *Buffer) flushLoop() {
// flushLocked acquires the lock, drains the buffer, and
// writes the data to a compressed file.
func (b *Buffer) flushLocked() error {
func (b *Buffer) flushLocked() {
b.mu.Lock()
if b.buf.Len() == 0 {
b.mu.Unlock()
return nil
return
}
data := b.drainBuf()
b.mu.Unlock()
return b.writeFile(data)
b.writeFile(data)
}
// drainBuf copies the buffer contents and resets it.
@@ -176,48 +164,42 @@ func (b *Buffer) drainBuf() []byte {
// writeFile creates a timestamped zstd-compressed JSONL file
// in the data directory.
func (b *Buffer) writeFile(data []byte) error {
func (b *Buffer) writeFile(data []byte) {
ts := time.Now().UTC().Format("2006-01-02T15-04-05.000Z")
name := fmt.Sprintf("reports-%s.jsonl.zst", ts)
path := filepath.Join(b.dataDir, name)
// path is built from the operator-supplied dataDir plus a
// generated timestamp, so it carries no external input.
f, err := os.OpenFile( //nolint:gosec // see comment above
f, err := os.OpenFile( //nolint:gosec // path built from controlled dataDir + timestamp
path,
os.O_WRONLY|os.O_CREATE|os.O_EXCL,
filePerms,
)
if err != nil {
return fmt.Errorf("create report file: %w", err)
b.log.Error("create report file", "error", err)
return
}
// Closes the file on the early returns below. The success
// path closes it explicitly to check the error; closing it
// a second time here is harmless.
defer func() { _ = f.Close() }()
enc, err := zstd.NewWriter(f)
if err != nil {
return fmt.Errorf("create zstd encoder: %w", err)
b.log.Error("create zstd encoder", "error", err)
return
}
_, err = enc.Write(data)
if err != nil {
_, writeErr := enc.Write(data)
if writeErr != nil {
b.log.Error("write compressed data", "error", writeErr)
_ = enc.Close()
return fmt.Errorf("write compressed data: %w", err)
return
}
err = enc.Close()
if err != nil {
return fmt.Errorf("close zstd encoder: %w", err)
closeErr := enc.Close()
if closeErr != nil {
b.log.Error("close zstd encoder", "error", closeErr)
}
err = f.Close()
if err != nil {
return fmt.Errorf("close report file: %w", err)
}
return nil
}
@@ -1,8 +1,6 @@
package reportbuf_test
import (
"errors"
"io/fs"
"os"
"strings"
"testing"
@@ -54,46 +52,6 @@ func TestFlushOnShutdown(t *testing.T) {
}
}
// TestFailedFinalFlushFailsStop proves a final flush that cannot
// write its file makes the stop fail, which makes the process
// exit non-zero instead of dropping the buffered reports silently.
func TestFailedFinalFlushFailsStop(t *testing.T) {
dir := t.TempDir()
t.Setenv("DATA_DIR", dir)
var buf *reportbuf.Buffer
app := fxtest.New(t,
fx.Provide(
globals.New,
logger.New,
config.New,
reportbuf.New,
),
fx.Populate(&buf),
)
app.RequireStart()
err := buf.Append(map[string]string{"probe": "shutdown"})
if err != nil {
t.Fatalf("append report: %v", err)
}
// Removing the data directory leaves the final flush nowhere to
// write. A read-only directory would not do: tests run as root
// in the backend image, and root ignores the read-only bit.
err = os.RemoveAll(dir)
if err != nil {
t.Fatalf("remove data dir: %v", err)
}
err = app.Stop(t.Context())
if !errors.Is(err, fs.ErrNotExist) {
t.Fatalf("stop error = %v, want the final flush's error", err)
}
}
func hasReportFile(t *testing.T, dir string) bool {
t.Helper()
-5
View File
@@ -1,5 +0,0 @@
package server
// MaxRequestBodyBytes exposes the router-wide body limit to the
// external tests.
const MaxRequestBodyBytes = maxRequestBodyBytes
+2 -10
View File
@@ -7,26 +7,18 @@ import (
"github.com/go-chi/chi/v5/middleware"
)
const (
requestTimeout = 60 * time.Second
// maxRequestBodyBytes caps every request body. A route group
// can mount s.mw.MaxBodyBytes with a smaller value to lower
// its bound, but cannot raise it: this cap runs first.
maxRequestBodyBytes int64 = 1 << 20 // 1 MiB
)
const requestTimeout = 60 * time.Second
// SetupRoutes configures the chi router with middleware and
// all application routes.
func (s *Server) SetupRoutes() {
s.router = chi.NewRouter()
s.router.Use(s.mw.Recoverer())
s.router.Use(middleware.Recoverer)
s.router.Use(middleware.RequestID)
s.router.Use(s.mw.Logging())
s.router.Use(s.mw.SecurityHeaders())
s.router.Use(s.mw.CORS())
s.router.Use(s.mw.MaxBodyBytes(maxRequestBodyBytes))
s.router.Use(middleware.Timeout(requestTimeout))
s.router.Get(
-76
View File
@@ -1,76 +0,0 @@
package server_test
import (
"net/http"
"net/http/httptest"
"strings"
"testing"
"sneak.berlin/go/netwatch/internal/config"
"sneak.berlin/go/netwatch/internal/globals"
"sneak.berlin/go/netwatch/internal/handlers"
"sneak.berlin/go/netwatch/internal/healthcheck"
"sneak.berlin/go/netwatch/internal/logger"
"sneak.berlin/go/netwatch/internal/middleware"
"sneak.berlin/go/netwatch/internal/reportbuf"
"sneak.berlin/go/netwatch/internal/server"
"go.uber.org/fx"
"go.uber.org/fx/fxtest"
)
// TestHealthCheckRejectsOversizeBody sends the health check, which
// never reads its body, a body one byte over the limit. Only the
// router-wide body limit can reject it.
func TestHealthCheckRejectsOversizeBody(t *testing.T) {
t.Parallel()
var srv *server.Server
// The same constructors as main, never started: SetupRoutes is
// called directly, so nothing listens.
app := fxtest.New(t,
fx.Provide(
config.New,
globals.New,
handlers.New,
healthcheck.New,
logger.New,
middleware.New,
reportbuf.New,
server.New,
),
fx.Populate(&srv),
)
err := app.Err()
if err != nil {
t.Fatalf("build server: %v", err)
}
srv.SetupRoutes()
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(),
http.MethodGet, "/.well-known/healthcheck",
strings.NewReader(
strings.Repeat("x", int(server.MaxRequestBodyBytes)+1),
),
)
srv.ServeHTTP(rec, req)
if rec.Code != http.StatusRequestEntityTooLarge {
t.Fatalf("status = %d, want %d",
rec.Code, http.StatusRequestEntityTooLarge)
}
if got := rec.Body.String(); got != "{\"status\":\"error\"}\n" {
t.Errorf("body = %q, want %q", got, "{\"status\":\"error\"}\n")
}
got := rec.Header().Get("Content-Type")
if got != "application/json; charset=utf-8" {
t.Errorf("Content-Type = %q, want a JSON content type", got)
}
}
+2 -14
View File
@@ -443,10 +443,6 @@ export function buildReport(hosts, clientId, now, since) {
// samples next interval (bounded by the history window — whatever has since
// fallen out is dropped). Failure is quiet: one debug line per outage, one
// on recovery, never an alert, never a tight retry loop.
//
// Only one report is ever in flight, and it is abandoned after half the
// interval, so a slow POST can neither overlap the next report (which would
// carry the same samples) nor stall reporting for good.
class Reporter {
constructor(state, clientId, intervalMs) {
this.state = state;
@@ -454,7 +450,6 @@ class Reporter {
this.intervalMs = intervalMs;
this.marks = new Map();
this.failing = false;
this.sending = false;
this.timerId = null;
}
@@ -464,7 +459,7 @@ class Reporter {
}
async flush() {
if (this.sending || this.state.paused) return;
if (this.state.paused) return;
const report = buildReport(
this.state.allHosts,
this.clientId,
@@ -472,20 +467,15 @@ class Reporter {
this.marks,
);
if (!report) return;
this.sending = true;
try {
const resp = await fetch("/api/v1/reports", {
method: "POST",
headers: { "Content-Type": "application/json" },
credentials: "omit",
body: JSON.stringify(report.body),
signal: AbortSignal.timeout(this.intervalMs / 2),
});
if (!resp.ok) throw new Error(`HTTP ${resp.status}`);
for (const [url, t] of report.marks) {
const current = this.marks.get(url) ?? 0;
this.marks.set(url, Math.max(current, t));
}
for (const [url, t] of report.marks) this.marks.set(url, t);
if (this.failing) {
log.debug("Report delivery recovered");
this.failing = false;
@@ -495,8 +485,6 @@ class Reporter {
log.debug(`Report delivery failed: ${err.message}`);
this.failing = true;
}
} finally {
this.sending = false;
}
}
}