Compare commits
3
Commits
02a503b9d7
...
ca59368081
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ca59368081 | ||
|
|
503399e020 | ||
|
|
7a1ee6e5a8 |
+2
-2
@@ -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.7.2 (2026-08-09)
|
||||
FROM golangci/golangci-lint@sha256:5d6d5c70a61f1356adfd9dd6316ce286799fefc9d743421356ff1b00842368ba AS lint
|
||||
# golangci/golangci-lint:v2.12.2 (2026-08-10)
|
||||
FROM golangci/golangci-lint@sha256:5cceeef04e53efe1470638d4b4b4f5ceefd574955ab3941b2d9a68a8c9ad5240 AS lint
|
||||
|
||||
WORKDIR /src
|
||||
COPY backend/go.mod backend/go.sum ./
|
||||
|
||||
@@ -23,6 +23,9 @@ docker build -t netwatch .
|
||||
docker run -p 8080:8080 netwatch
|
||||
```
|
||||
|
||||
`yarn dev` proxies `/api` to `http://127.0.0.1:8080`, so a locally running
|
||||
`netwatch-server` (see `backend/`) receives the reports the page posts.
|
||||
|
||||
## Entrypoints
|
||||
|
||||
This repository adheres to the
|
||||
@@ -85,6 +88,18 @@ code lives in `src/main.js` with a class-based architecture:
|
||||
- **`tick()`**: Main loop — measures all hosts in parallel via `Promise.all`,
|
||||
pushes samples, redraws UI. When paused, pushes blank markers (no probes, no
|
||||
false outage)
|
||||
- **`Reporter`**: Posts collected samples to the backend
|
||||
|
||||
### Reporting
|
||||
|
||||
Every `reportInterval` (default 60s) the page POSTs a JSON report to the
|
||||
same-origin path `/api/v1/reports`: a random per-browser `clientId` kept in
|
||||
`localStorage`, `geo` sent as null, and each host's unreported, non-paused
|
||||
samples (timestamp, latency, error). A per-host high-water mark makes every
|
||||
report a delta, so only new samples are sent; the mark advances only on a
|
||||
delivered report, and while paused nothing is sent. Delivery failure is quiet —
|
||||
one debug-log line per outage, retried at the next interval, never blocking
|
||||
probing. The report-building step is a pure function of host state.
|
||||
|
||||
### Monitoring targets
|
||||
|
||||
|
||||
@@ -10,18 +10,29 @@
|
||||
|
||||
# Status
|
||||
|
||||
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.
|
||||
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.
|
||||
|
||||
# Next Step
|
||||
|
||||
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.
|
||||
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.
|
||||
|
||||
# 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
|
||||
@@ -31,11 +42,26 @@ files, so merging it also closes most compliance gaps.
|
||||
`OnStop` is idempotent; and `writeTimeout` now exceeds the chi per-request
|
||||
budget so that budget is actually reachable. Dead `startupTime`, `exitCode`,
|
||||
and `cancelFunc` fields were removed
|
||||
- 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`
|
||||
- 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
|
||||
@@ -65,7 +91,7 @@ files, so merging it also closes most compliance gaps.
|
||||
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, unmerged)
|
||||
golangci-lint) and moved to repo root (feat/reportbuf-storage)
|
||||
- 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,
|
||||
@@ -87,3 +113,9 @@ files, so merging it also closes most compliance gaps.
|
||||
(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>
|
||||
|
||||
+14
-12
@@ -1,5 +1,9 @@
|
||||
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
|
||||
@@ -14,19 +18,17 @@ linters:
|
||||
- wsl # Deprecated, replaced by wsl_v5
|
||||
- wrapcheck # Too verbose for internal packages
|
||||
- varnamelen # Short names like db, id are idiomatic Go
|
||||
|
||||
linters-settings:
|
||||
lll:
|
||||
line-length: 88
|
||||
funlen:
|
||||
lines: 80
|
||||
statements: 50
|
||||
cyclop:
|
||||
max-complexity: 15
|
||||
dupl:
|
||||
threshold: 100
|
||||
settings:
|
||||
lll:
|
||||
line-length: 88
|
||||
funlen:
|
||||
lines: 80
|
||||
statements: 50
|
||||
cyclop:
|
||||
max-complexity: 15
|
||||
dupl:
|
||||
threshold: 100
|
||||
|
||||
issues:
|
||||
exclude-use-default: false
|
||||
max-issues-per-linter: 0
|
||||
max-same-issues: 0
|
||||
|
||||
@@ -10,6 +10,17 @@ 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
|
||||
@@ -22,6 +33,14 @@ 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:
|
||||
|
||||
@@ -0,0 +1,13 @@
|
||||
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}
|
||||
}
|
||||
@@ -18,6 +18,13 @@ 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
|
||||
@@ -30,7 +37,7 @@ type Params struct {
|
||||
|
||||
// Handlers provides HTTP handler factories for all endpoints.
|
||||
type Handlers struct {
|
||||
buf *reportbuf.Buffer
|
||||
buf reportAppender
|
||||
hc *healthcheck.Healthcheck
|
||||
log *slog.Logger
|
||||
params *Params
|
||||
@@ -72,3 +79,15 @@ 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)
|
||||
}
|
||||
|
||||
@@ -2,10 +2,14 @@ package handlers
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"net/http"
|
||||
)
|
||||
|
||||
const maxReportBodyBytes = 1 << 20 // 1 MiB
|
||||
// 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
|
||||
|
||||
type reportSample struct {
|
||||
T int64 `json:"t"`
|
||||
@@ -35,48 +39,80 @@ 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 := json.NewDecoder(r.Body).Decode(&rpt)
|
||||
err := s.decodeJSON(w, r, &rpt)
|
||||
if err != nil {
|
||||
s.log.Error("failed to decode report",
|
||||
"error", err,
|
||||
)
|
||||
s.respondJSON(w, r,
|
||||
&response{Status: "error"},
|
||||
http.StatusBadRequest,
|
||||
s.decodeErrorStatus(err),
|
||||
)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
totalSamples := 0
|
||||
for _, h := range rpt.Hosts {
|
||||
totalSamples += len(h.History)
|
||||
}
|
||||
s.logReportReceived(rpt)
|
||||
|
||||
s.log.Info("report received",
|
||||
"client_id", rpt.ClientID,
|
||||
"timestamp", rpt.Timestamp,
|
||||
"host_count", len(rpt.Hosts),
|
||||
"total_samples", totalSamples,
|
||||
"geo", string(rpt.Geo),
|
||||
)
|
||||
|
||||
bufErr := s.buf.Append(rpt)
|
||||
if bufErr != nil {
|
||||
s.log.Error("failed to buffer report",
|
||||
"error", bufErr,
|
||||
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,
|
||||
)
|
||||
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()),
|
||||
)
|
||||
|
||||
return http.StatusBadRequest
|
||||
}
|
||||
|
||||
// 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),
|
||||
"host_count", len(rpt.Hosts),
|
||||
"total_samples", totalSamples,
|
||||
"geo_bytes", len(rpt.Geo),
|
||||
)
|
||||
}
|
||||
|
||||
// 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]
|
||||
}
|
||||
|
||||
return s
|
||||
}
|
||||
|
||||
@@ -0,0 +1,212 @@
|
||||
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,6 +1,7 @@
|
||||
package middleware
|
||||
|
||||
import (
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/netip"
|
||||
)
|
||||
@@ -8,6 +9,12 @@ 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,
|
||||
|
||||
@@ -3,11 +3,14 @@
|
||||
package middleware
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/netip"
|
||||
"runtime/debug"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
@@ -22,6 +25,14 @@ 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.
|
||||
@@ -236,6 +247,79 @@ 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{
|
||||
|
||||
@@ -1,14 +1,28 @@
|
||||
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()
|
||||
|
||||
@@ -41,32 +55,32 @@ func clientIPCases() []clientIPCase {
|
||||
return []clientIPCase{
|
||||
{
|
||||
name: "trusted proxy uses forwarded-for",
|
||||
remoteAddr: "127.0.0.1:5000",
|
||||
xff: "203.0.113.7",
|
||||
want: "203.0.113.7",
|
||||
remoteAddr: loopbackPeer,
|
||||
xff: forwardedIP,
|
||||
want: forwardedIP,
|
||||
},
|
||||
{
|
||||
name: "trusted proxy uses left-most of chain",
|
||||
remoteAddr: "10.1.2.3:5000",
|
||||
xff: "203.0.113.7, 10.1.2.3",
|
||||
want: "203.0.113.7",
|
||||
xff: forwardedIP + ", 10.1.2.3",
|
||||
want: forwardedIP,
|
||||
},
|
||||
{
|
||||
name: "trusted proxy falls back to x-real-ip",
|
||||
remoteAddr: "127.0.0.1:5000",
|
||||
xRealIP: "203.0.113.9",
|
||||
want: "203.0.113.9",
|
||||
remoteAddr: loopbackPeer,
|
||||
xRealIP: realIP,
|
||||
want: realIP,
|
||||
},
|
||||
{
|
||||
name: "untrusted peer ignores forwarded-for",
|
||||
remoteAddr: "198.51.100.4:5000",
|
||||
xff: "203.0.113.7",
|
||||
xff: forwardedIP,
|
||||
want: "198.51.100.4",
|
||||
},
|
||||
{
|
||||
name: "untrusted peer ignores x-real-ip",
|
||||
remoteAddr: "198.51.100.4:5000",
|
||||
xRealIP: "203.0.113.9",
|
||||
xRealIP: realIP,
|
||||
want: "198.51.100.4",
|
||||
},
|
||||
{
|
||||
@@ -76,7 +90,7 @@ func clientIPCases() []clientIPCase {
|
||||
},
|
||||
{
|
||||
name: "trusted proxy with garbage header uses peer",
|
||||
remoteAddr: "127.0.0.1:5000",
|
||||
remoteAddr: loopbackPeer,
|
||||
xff: "not-an-ip",
|
||||
want: "127.0.0.1",
|
||||
},
|
||||
@@ -119,7 +133,7 @@ func TestSecurityHeaders(t *testing.T) {
|
||||
)
|
||||
|
||||
rec := httptest.NewRecorder()
|
||||
req := httptest.NewRequest(http.MethodGet, "/", http.NoBody)
|
||||
req := httptest.NewRequestWithContext(t.Context(), http.MethodGet, "/", http.NoBody)
|
||||
handler.ServeHTTP(rec, req)
|
||||
|
||||
want := map[string]string{
|
||||
@@ -137,3 +151,152 @@ 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())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -80,12 +80,16 @@ 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)
|
||||
b.flushLocked()
|
||||
err = b.flushLocked()
|
||||
})
|
||||
|
||||
return nil
|
||||
// A failed final flush fails the stop, so the process
|
||||
// exits non-zero.
|
||||
return err
|
||||
},
|
||||
})
|
||||
|
||||
@@ -109,7 +113,12 @@ func (b *Buffer) Append(v any) error {
|
||||
data := b.drainBuf()
|
||||
b.mu.Unlock()
|
||||
|
||||
go b.writeFile(data)
|
||||
go func() {
|
||||
writeErr := b.writeFile(data)
|
||||
if writeErr != nil {
|
||||
b.log.Error("flush reports failed", "error", writeErr)
|
||||
}
|
||||
}()
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -128,7 +137,10 @@ func (b *Buffer) flushLoop() {
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
b.flushLocked()
|
||||
err := b.flushLocked()
|
||||
if err != nil {
|
||||
b.log.Error("flush reports failed", "error", err)
|
||||
}
|
||||
case <-b.done:
|
||||
return
|
||||
}
|
||||
@@ -137,19 +149,19 @@ func (b *Buffer) flushLoop() {
|
||||
|
||||
// flushLocked acquires the lock, drains the buffer, and
|
||||
// writes the data to a compressed file.
|
||||
func (b *Buffer) flushLocked() {
|
||||
func (b *Buffer) flushLocked() error {
|
||||
b.mu.Lock()
|
||||
|
||||
if b.buf.Len() == 0 {
|
||||
b.mu.Unlock()
|
||||
|
||||
return
|
||||
return nil
|
||||
}
|
||||
|
||||
data := b.drainBuf()
|
||||
b.mu.Unlock()
|
||||
|
||||
b.writeFile(data)
|
||||
return b.writeFile(data)
|
||||
}
|
||||
|
||||
// drainBuf copies the buffer contents and resets it.
|
||||
@@ -164,42 +176,48 @@ func (b *Buffer) drainBuf() []byte {
|
||||
|
||||
// writeFile creates a timestamped zstd-compressed JSONL file
|
||||
// in the data directory.
|
||||
func (b *Buffer) writeFile(data []byte) {
|
||||
func (b *Buffer) writeFile(data []byte) error {
|
||||
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)
|
||||
|
||||
f, err := os.OpenFile( //nolint:gosec // path built from controlled dataDir + timestamp
|
||||
// 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
|
||||
path,
|
||||
os.O_WRONLY|os.O_CREATE|os.O_EXCL,
|
||||
filePerms,
|
||||
)
|
||||
if err != nil {
|
||||
b.log.Error("create report file", "error", err)
|
||||
|
||||
return
|
||||
return fmt.Errorf("create report file: %w", err)
|
||||
}
|
||||
|
||||
// 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 {
|
||||
b.log.Error("create zstd encoder", "error", err)
|
||||
|
||||
return
|
||||
return fmt.Errorf("create zstd encoder: %w", err)
|
||||
}
|
||||
|
||||
_, writeErr := enc.Write(data)
|
||||
if writeErr != nil {
|
||||
b.log.Error("write compressed data", "error", writeErr)
|
||||
|
||||
_, err = enc.Write(data)
|
||||
if err != nil {
|
||||
_ = enc.Close()
|
||||
|
||||
return
|
||||
return fmt.Errorf("write compressed data: %w", err)
|
||||
}
|
||||
|
||||
closeErr := enc.Close()
|
||||
if closeErr != nil {
|
||||
b.log.Error("close zstd encoder", "error", closeErr)
|
||||
err = enc.Close()
|
||||
if err != nil {
|
||||
return fmt.Errorf("close zstd encoder: %w", err)
|
||||
}
|
||||
|
||||
err = f.Close()
|
||||
if err != nil {
|
||||
return fmt.Errorf("close report file: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
package reportbuf_test
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"io/fs"
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
@@ -52,6 +54,46 @@ 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()
|
||||
|
||||
|
||||
@@ -0,0 +1,5 @@
|
||||
package server
|
||||
|
||||
// MaxRequestBodyBytes exposes the router-wide body limit to the
|
||||
// external tests.
|
||||
const MaxRequestBodyBytes = maxRequestBodyBytes
|
||||
@@ -7,18 +7,26 @@ import (
|
||||
"github.com/go-chi/chi/v5/middleware"
|
||||
)
|
||||
|
||||
const requestTimeout = 60 * time.Second
|
||||
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
|
||||
)
|
||||
|
||||
// SetupRoutes configures the chi router with middleware and
|
||||
// all application routes.
|
||||
func (s *Server) SetupRoutes() {
|
||||
s.router = chi.NewRouter()
|
||||
|
||||
s.router.Use(middleware.Recoverer)
|
||||
s.router.Use(s.mw.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(
|
||||
|
||||
@@ -0,0 +1,76 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
+177
-4
@@ -7,9 +7,11 @@ import "./styles.css";
|
||||
// graphMaxLatency — values above it pin to the top of the chart but still
|
||||
// display their real value in the latency figure. The history buffer holds
|
||||
// maxHistoryPoints samples (historyDuration / updateInterval).
|
||||
// reportInterval is how often collected samples are POSTed to the backend.
|
||||
const CONFIG = {
|
||||
updateInterval: 3000,
|
||||
maxHistoryPoints: 100,
|
||||
reportInterval: 60000,
|
||||
get historyDuration() {
|
||||
return (this.maxHistoryPoints * this.updateInterval) / 1000;
|
||||
},
|
||||
@@ -346,6 +348,159 @@ class AppState {
|
||||
}
|
||||
}
|
||||
|
||||
// --- Reporting ---------------------------------------------------------------
|
||||
|
||||
// A random UUIDv4. `crypto.randomUUID` exists only in secure contexts
|
||||
// (HTTPS or localhost); over plain HTTP to any other host — the normal LAN
|
||||
// deployment — it is undefined, so feature-detect it and otherwise build the
|
||||
// id from `crypto.getRandomValues`, which is available in insecure contexts.
|
||||
function randomId() {
|
||||
if (typeof crypto !== "undefined" && crypto.randomUUID) {
|
||||
return crypto.randomUUID();
|
||||
}
|
||||
const bytes = new Uint8Array(16);
|
||||
crypto.getRandomValues(bytes);
|
||||
bytes[6] = (bytes[6] & 0x0f) | 0x40; // version 4
|
||||
bytes[8] = (bytes[8] & 0x3f) | 0x80; // variant 1
|
||||
const hex = [...bytes].map((b) => b.toString(16).padStart(2, "0"));
|
||||
return (
|
||||
hex.slice(0, 4).join("") +
|
||||
"-" +
|
||||
hex.slice(4, 6).join("") +
|
||||
"-" +
|
||||
hex.slice(6, 8).join("") +
|
||||
"-" +
|
||||
hex.slice(8, 10).join("") +
|
||||
"-" +
|
||||
hex.slice(10, 16).join("")
|
||||
);
|
||||
}
|
||||
|
||||
// A random id identifying this browser across reports. Generated once and
|
||||
// kept in localStorage; if storage is unavailable (e.g. private mode) a
|
||||
// fresh id is used for this session only.
|
||||
function getClientId() {
|
||||
const key = "netwatch-client-id";
|
||||
try {
|
||||
let id = localStorage.getItem(key);
|
||||
if (!id) {
|
||||
id = randomId();
|
||||
localStorage.setItem(key, id);
|
||||
}
|
||||
return id;
|
||||
} catch {
|
||||
return randomId();
|
||||
}
|
||||
}
|
||||
|
||||
// Build the delta report body the backend decodes, plus the new per-host
|
||||
// high-water marks. Pure function of the passed state: `hosts` is an array
|
||||
// of { name, url, status, history }, `since` maps a host url to the Unix-ms
|
||||
// timestamp of the last sample already reported for it, and `now` is a Date.
|
||||
// Only non-paused samples newer than the mark are included. Returns null
|
||||
// when no host has an unreported sample.
|
||||
export function buildReport(hosts, clientId, now, since) {
|
||||
const reportHosts = [];
|
||||
const marks = new Map();
|
||||
for (const host of hosts) {
|
||||
const mark = since.get(host.url) ?? 0;
|
||||
const samples = [];
|
||||
let high = mark;
|
||||
for (const p of host.history) {
|
||||
if (p.paused) continue;
|
||||
if (p.timestamp <= mark) continue;
|
||||
samples.push({
|
||||
t: p.timestamp,
|
||||
latency: p.latency,
|
||||
error: p.error ?? null,
|
||||
});
|
||||
if (p.timestamp > high) high = p.timestamp;
|
||||
}
|
||||
if (samples.length === 0) continue;
|
||||
reportHosts.push({
|
||||
name: host.name,
|
||||
url: host.url,
|
||||
status: host.status,
|
||||
history: samples,
|
||||
});
|
||||
marks.set(host.url, high);
|
||||
}
|
||||
if (reportHosts.length === 0) return null;
|
||||
return {
|
||||
body: {
|
||||
clientId,
|
||||
geo: null,
|
||||
hosts: reportHosts,
|
||||
timestamp: now.toISOString(),
|
||||
},
|
||||
marks,
|
||||
};
|
||||
}
|
||||
|
||||
// Periodically POSTs unreported samples to the same-origin backend. Holds
|
||||
// the per-host high-water marks so each report is a delta; marks only
|
||||
// advance on a delivered report, so a failed POST simply re-sends those
|
||||
// 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;
|
||||
this.clientId = clientId;
|
||||
this.intervalMs = intervalMs;
|
||||
this.marks = new Map();
|
||||
this.failing = false;
|
||||
this.sending = false;
|
||||
this.timerId = null;
|
||||
}
|
||||
|
||||
start() {
|
||||
if (this.timerId) return;
|
||||
this.timerId = setInterval(() => this.flush(), this.intervalMs);
|
||||
}
|
||||
|
||||
async flush() {
|
||||
if (this.sending || this.state.paused) return;
|
||||
const report = buildReport(
|
||||
this.state.allHosts,
|
||||
this.clientId,
|
||||
new Date(),
|
||||
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));
|
||||
}
|
||||
if (this.failing) {
|
||||
log.debug("Report delivery recovered");
|
||||
this.failing = false;
|
||||
}
|
||||
} catch (err) {
|
||||
if (!this.failing) {
|
||||
log.debug(`Report delivery failed: ${err.message}`);
|
||||
this.failing = true;
|
||||
}
|
||||
} finally {
|
||||
this.sending = false;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// --- Latency Measurement -----------------------------------------------------
|
||||
|
||||
async function measureLatency(url) {
|
||||
@@ -1162,6 +1317,19 @@ async function init() {
|
||||
buildUI(state);
|
||||
log.info("UI built, starting tick loop");
|
||||
|
||||
// Reporting is best-effort: any failure setting it up (e.g. no usable
|
||||
// crypto for the client id) must never stop the monitor from probing.
|
||||
try {
|
||||
const reporter = new Reporter(
|
||||
state,
|
||||
getClientId(),
|
||||
CONFIG.reportInterval,
|
||||
);
|
||||
reporter.start();
|
||||
} catch (err) {
|
||||
log.error(`Reporting disabled: ${err.message}`);
|
||||
}
|
||||
|
||||
document
|
||||
.getElementById("pause-btn")
|
||||
.addEventListener("click", () => togglePause(state));
|
||||
@@ -1274,8 +1442,13 @@ async function init() {
|
||||
setTimeout(() => handleResize(state), 100);
|
||||
}
|
||||
|
||||
if (document.readyState === "loading") {
|
||||
document.addEventListener("DOMContentLoaded", init);
|
||||
} else {
|
||||
init();
|
||||
// Bootstrap only when loaded as the page: a real DOM containing the #app
|
||||
// mount point this module renders into. Importing the module in a unit test
|
||||
// (which has no #app) runs nothing, so buildReport can be tested in isolation.
|
||||
if (typeof document !== "undefined" && document.getElementById("app")) {
|
||||
if (document.readyState === "loading") {
|
||||
document.addEventListener("DOMContentLoaded", init);
|
||||
} else {
|
||||
init();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,6 +7,13 @@ const commitFull = execSync("git rev-parse HEAD").toString().trim();
|
||||
|
||||
export default defineConfig({
|
||||
plugins: [tailwindcss()],
|
||||
server: {
|
||||
// Proxy /api to a locally running netwatch-server so `yarn dev`
|
||||
// exercises the real report-posting path.
|
||||
proxy: {
|
||||
"/api": "http://127.0.0.1:8080",
|
||||
},
|
||||
},
|
||||
define: {
|
||||
__COMMIT_HASH__: JSON.stringify(commitHash),
|
||||
__COMMIT_FULL__: JSON.stringify(commitFull),
|
||||
|
||||
Reference in New Issue
Block a user