2 Commits
Author SHA1 Message Date
sneak 924d4b7177 fix(server): shut down through fx so buffered reports flush (closes #22)
check / check (push) Failing after 1s
The server ran os.Exit at the end of its own goroutine, racing fx's
teardown and sometimes killing the process before reportbuf's OnStop
flushed — silently losing a full flush window of telemetry on every
restart, at exit 0. Shutdown now goes through fx.Shutdowner, so every
OnStop runs in order.

The http.Server is built synchronously in OnStart before the serving
goroutine, so shutdown can no longer race or nil-deref it. A listen
failure exits non-zero via fx.ExitCode(1). reportbuf's OnStop is guarded
by sync.Once. writeTimeout now exceeds the chi per-request budget so that
budget is reachable. Dead startupTime, exitCode, and cancelFunc fields
are gone. A new test asserts a buffered report reaches disk after the
lifecycle stops.

Model: opus-4-8
2026-09-21 13:10:22 +00:00
clawbot f3895789d2 feat(backend): server hardening: timeouts, security headers, trusted-proxy client IP (closes #19)
check / check (push) Failing after 1s
Add ReadHeaderTimeout and IdleTimeout to the http.Server as named constants beside the existing timeouts. Add a SecurityHeaders middleware (HSTS, a JSON-API CSP of default-src 'none'; frame-ancestors 'none', X-Frame-Options DENY, nosniff, Referrer-Policy, Permissions-Policy), registered before CORS so preflight responses carry it. Resolve the client IP from X-Forwarded-For / X-Real-IP only when the direct peer is in the trusted-proxy allowlist (loopback plus RFC1918 by default, configurable via TRUSTED_PROXIES); an untrusted peer's forwarded headers are ignored. Uses net/netip; no new dependency.

Model: opus-4-8 (implementation and review); claude-fable-5 (merge)
2026-09-21 15:05:25 +02:00
11 changed files with 500 additions and 115 deletions
+14
View File
@@ -22,6 +22,20 @@ files, so merging it also closes most compliance gaps.
# Completed Steps
- 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
telemetry was silently lost on every restart. The `http.Server` is now built
before its serving goroutine starts, so shutdown can no longer race or
nil-deref it; a listen failure exits non-zero via `fx.Shutdowner`; `reportbuf`
`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: 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: 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
+7 -1
View File
@@ -43,10 +43,16 @@ Internal packages in `internal/` follow standard Go project layout:
### Configuration
| Variable | Default | Description |
| ---------- | ------------------ | --------------------------------- |
| ----------------- | -------------------- | -------------------------------------------------------------------------------------------------------- |
| `PORT` | `8080` | HTTP listen port |
| `DATA_DIR` | `./data/reports` | Directory for compressed reports |
| `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 |
`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
RFC1918 ranges match `nginx.conf`. A request whose direct peer is outside this
set has its forwarded headers ignored, and the direct peer is logged instead.
### Report storage
+28
View File
@@ -5,6 +5,7 @@ package config
import (
"errors"
"log/slog"
"strings"
"sneak.berlin/go/netwatch/internal/globals"
"sneak.berlin/go/netwatch/internal/logger"
@@ -14,6 +15,14 @@ import (
"go.uber.org/fx"
)
// defaultTrustedProxies lists the networks whose forwarded
// headers are honoured by default. It covers the RFC1918
// ranges (to match nginx.conf) plus IPv4 and IPv6 loopback,
// because the reverse proxy shares the container and reaches
// the backend over loopback.
const defaultTrustedProxies = "127.0.0.1/32,::1/128," +
"10.0.0.0/8,172.16.0.0/12,192.168.0.0/16"
// Params defines the dependencies for Config.
type Params struct {
fx.In
@@ -30,6 +39,7 @@ type Config struct {
MetricsUsername string
Port int
SentryDSN string
TrustedProxies []string
log *slog.Logger
params *Params
}
@@ -56,6 +66,7 @@ func New(
viper.SetDefault("SENTRY_DSN", "")
viper.SetDefault("METRICS_USERNAME", "")
viper.SetDefault("METRICS_PASSWORD", "")
viper.SetDefault("TRUSTED_PROXIES", defaultTrustedProxies)
err := viper.ReadInConfig()
if err != nil {
@@ -73,6 +84,7 @@ func New(
MetricsUsername: viper.GetString("METRICS_USERNAME"),
Port: viper.GetInt("PORT"),
SentryDSN: viper.GetString("SENTRY_DSN"),
TrustedProxies: splitList(viper.GetString("TRUSTED_PROXIES")),
log: log,
params: &params,
}
@@ -84,3 +96,19 @@ func New(
return s, nil
}
// splitList turns a comma-separated setting into a trimmed
// slice, dropping empty entries.
func splitList(raw string) []string {
parts := strings.Split(raw, ",")
out := make([]string, 0, len(parts))
for _, p := range parts {
p = strings.TrimSpace(p)
if p != "" {
out = append(out, p)
}
}
return out
}
@@ -0,0 +1,21 @@
package middleware
import (
"net/http"
"net/netip"
)
// Test-only wrappers exposing unexported helpers to the
// external middleware_test package.
func ClientIP(
remoteAddr string,
header http.Header,
trusted []netip.Prefix,
) string {
return clientIP(remoteAddr, header, trusted)
}
func ParseTrustedProxies(cidrs []string) ([]netip.Prefix, error) {
return parseTrustedProxies(cidrs)
}
+128 -1
View File
@@ -3,9 +3,12 @@
package middleware
import (
"fmt"
"log/slog"
"net"
"net/http"
"net/netip"
"strings"
"time"
"sneak.berlin/go/netwatch/internal/config"
@@ -19,6 +22,15 @@ import (
const corsMaxAgeSec = 300
// Security header values. The backend is a JSON API with no
// HTML surface, so the CSP forbids every resource type and
// framing outright.
const (
hstsValue = "max-age=31536000; includeSubDomains"
cspValue = "default-src 'none'; frame-ancestors 'none'"
permissionsPolicyValue = "camera=(), microphone=(), geolocation=()"
)
// Params defines the dependencies for Middleware.
type Params struct {
fx.In
@@ -32,6 +44,7 @@ type Params struct {
type Middleware struct {
log *slog.Logger
params *Params
trustedProxies []netip.Prefix
}
// New creates a Middleware instance.
@@ -39,13 +52,38 @@ func New(
_ fx.Lifecycle,
params Params,
) (*Middleware, error) {
trusted, err := parseTrustedProxies(params.Config.TrustedProxies)
if err != nil {
return nil, err
}
s := new(Middleware)
s.params = &params
s.log = params.Logger.Get()
s.trustedProxies = trusted
return s, nil
}
// parseTrustedProxies converts CIDR strings into prefixes,
// failing fast on any malformed entry.
func parseTrustedProxies(cidrs []string) ([]netip.Prefix, error) {
prefixes := make([]netip.Prefix, 0, len(cidrs))
for _, cidr := range cidrs {
prefix, err := netip.ParsePrefix(cidr)
if err != nil {
return nil, fmt.Errorf(
"trusted proxy %q: %w", cidr, err,
)
}
prefixes = append(prefixes, prefix.Masked())
}
return prefixes, nil
}
type loggingResponseWriter struct {
http.ResponseWriter
@@ -72,6 +110,70 @@ func ipFromHostPort(hostPort string) string {
return host
}
// clientIP resolves the caller's address. X-Forwarded-For and
// X-Real-IP are honoured only when the direct peer is a
// trusted proxy; otherwise the direct peer is returned so a
// spoofed header cannot forge the logged address.
func clientIP(
remoteAddr string,
header http.Header,
trusted []netip.Prefix,
) string {
peer := ipFromHostPort(remoteAddr)
if !addrInAny(peer, trusted) {
return peer
}
if xff := firstForwardedFor(header.Get("X-Forwarded-For")); xff != "" {
return xff
}
if xr := strings.TrimSpace(header.Get("X-Real-IP")); validIP(xr) {
return xr
}
return peer
}
// firstForwardedFor returns the left-most valid address in an
// X-Forwarded-For list (the original client), or "" if none.
func firstForwardedFor(value string) string {
for part := range strings.SplitSeq(value, ",") {
candidate := strings.TrimSpace(part)
if validIP(candidate) {
return candidate
}
}
return ""
}
func validIP(s string) bool {
_, err := netip.ParseAddr(s)
return err == nil
}
// addrInAny reports whether s parses as an address contained
// in any of the trusted prefixes.
func addrInAny(s string, trusted []netip.Prefix) bool {
addr, err := netip.ParseAddr(s)
if err != nil {
return false
}
addr = addr.Unmap()
for _, prefix := range trusted {
if prefix.Contains(addr) {
return true
}
}
return false
}
// Logging returns middleware that logs each request with
// timing, status code, and client information.
func (s *Middleware) Logging() func(http.Handler) http.Handler {
@@ -96,7 +198,11 @@ func (s *Middleware) Logging() func(http.Handler) http.Handler {
"referer", r.Referer(),
"proto", r.Proto,
"remote_ip",
ipFromHostPort(r.RemoteAddr),
clientIP(
r.RemoteAddr,
r.Header,
s.trustedProxies,
),
"status", lrw.statusCode,
"latency_ms",
latency.Milliseconds(),
@@ -109,6 +215,27 @@ func (s *Middleware) Logging() func(http.Handler) http.Handler {
}
}
// SecurityHeaders returns middleware that sets response
// security headers. It runs before CORS so the headers are
// present on preflight responses the CORS handler writes.
func (s *Middleware) SecurityHeaders() func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
h := w.Header()
h.Set("Strict-Transport-Security", hstsValue)
h.Set("Content-Security-Policy", cspValue)
h.Set("X-Frame-Options", "DENY")
h.Set("X-Content-Type-Options", "nosniff")
h.Set("Referrer-Policy", "no-referrer")
h.Set("Permissions-Policy", permissionsPolicyValue)
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{
@@ -0,0 +1,139 @@
package middleware_test
import (
"net/http"
"net/http/httptest"
"net/netip"
"testing"
"sneak.berlin/go/netwatch/internal/middleware"
)
func mustPrefixes(t *testing.T, cidrs ...string) []netip.Prefix {
t.Helper()
prefixes, err := middleware.ParseTrustedProxies(cidrs)
if err != nil {
t.Fatalf("ParseTrustedProxies(%v): %v", cidrs, err)
}
return prefixes
}
func TestParseTrustedProxiesRejectsMalformed(t *testing.T) {
t.Parallel()
_, err := middleware.ParseTrustedProxies([]string{"not-a-cidr"})
if err == nil {
t.Fatal("expected error for malformed CIDR, got nil")
}
}
type clientIPCase struct {
name string
remoteAddr string
xff string
xRealIP string
want string
}
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",
},
{
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",
},
{
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",
},
{
name: "untrusted peer ignores forwarded-for",
remoteAddr: "198.51.100.4:5000",
xff: "203.0.113.7",
want: "198.51.100.4",
},
{
name: "untrusted peer ignores x-real-ip",
remoteAddr: "198.51.100.4:5000",
xRealIP: "203.0.113.9",
want: "198.51.100.4",
},
{
name: "trusted proxy with no headers uses peer",
remoteAddr: "10.1.2.3:5000",
want: "10.1.2.3",
},
{
name: "trusted proxy with garbage header uses peer",
remoteAddr: "127.0.0.1:5000",
xff: "not-an-ip",
want: "127.0.0.1",
},
}
}
func TestClientIP(t *testing.T) {
t.Parallel()
trusted := mustPrefixes(t, "127.0.0.1/32", "::1/128", "10.0.0.0/8")
for _, tc := range clientIPCases() {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
header := http.Header{}
if tc.xff != "" {
header.Set("X-Forwarded-For", tc.xff)
}
if tc.xRealIP != "" {
header.Set("X-Real-IP", tc.xRealIP)
}
got := middleware.ClientIP(tc.remoteAddr, header, trusted)
if got != tc.want {
t.Errorf("ClientIP() = %q, want %q", got, tc.want)
}
})
}
}
func TestSecurityHeaders(t *testing.T) {
t.Parallel()
handler := (&middleware.Middleware{}).SecurityHeaders()(
http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
}),
)
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/", http.NoBody)
handler.ServeHTTP(rec, req)
want := map[string]string{
"Strict-Transport-Security": "max-age=31536000; includeSubDomains",
"Content-Security-Policy": "default-src 'none'; frame-ancestors 'none'",
"X-Frame-Options": "DENY",
"X-Content-Type-Options": "nosniff",
"Referrer-Policy": "no-referrer",
"Permissions-Policy": "camera=(), microphone=(), geolocation=()",
}
for name, value := range want {
if got := rec.Header().Get(name); got != value {
t.Errorf("header %s = %q, want %q", name, got, value)
}
}
}
+6
View File
@@ -45,6 +45,7 @@ type Buffer struct {
done chan struct{}
log *slog.Logger
mu sync.Mutex
stopOnce sync.Once
}
// New creates a Buffer and registers lifecycle hooks to
@@ -76,8 +77,13 @@ func New(
return nil
},
OnStop: func(_ context.Context) error {
// stopOnce makes OnStop idempotent: a second
// invocation must not close an already-closed channel
// (which would panic) or flush again.
b.stopOnce.Do(func() {
close(b.done)
b.flushLocked()
})
return nil
},
+69 -5
View File
@@ -1,13 +1,77 @@
package reportbuf_test
import (
"os"
"strings"
"testing"
_ "sneak.berlin/go/netwatch/internal/reportbuf"
"sneak.berlin/go/netwatch/internal/config"
"sneak.berlin/go/netwatch/internal/globals"
"sneak.berlin/go/netwatch/internal/logger"
"sneak.berlin/go/netwatch/internal/reportbuf"
"go.uber.org/fx"
"go.uber.org/fx/fxtest"
)
func TestImport(t *testing.T) {
t.Parallel()
// Compilation check — verifies the package parses
// and all imports resolve.
// TestFlushOnShutdown proves the flush-on-shutdown path: a
// report appended after start but before the periodic flush
// window must reach disk when the fx lifecycle stops. This is
// the exact case that silent data loss on restart used to
// destroy.
func TestFlushOnShutdown(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)
}
// RequireStop runs the reportbuf OnStop hook, which is the
// only code path that flushes buffered reports on shutdown.
app.RequireStop()
if !hasReportFile(t, dir) {
t.Fatal("no report file on disk after shutdown; " +
"the buffered report was lost")
}
}
func hasReportFile(t *testing.T, dir string) bool {
t.Helper()
entries, err := os.ReadDir(dir)
if err != nil {
t.Fatalf("read data dir: %v", err)
}
for _, e := range entries {
if strings.HasSuffix(e.Name(), ".jsonl.zst") {
info, statErr := e.Info()
if statErr != nil {
t.Fatalf("stat %s: %v", e.Name(), statErr)
}
if info.Size() > 0 {
return true
}
}
}
return false
}
+33 -10
View File
@@ -5,39 +5,62 @@ import (
"fmt"
"net/http"
"time"
"go.uber.org/fx"
)
const (
readTimeout = 10 * time.Second
writeTimeout = 10 * time.Second
readHeaderTimeout = 5 * time.Second
idleTimeout = 60 * time.Second
maxHeaderBytes = 1 << 20 // 1 MiB
// requestTimeout (routes.go) is the single per-request
// processing budget, enforced by chi's middleware.Timeout.
// writeTimeout must exceed that budget so a handler can write
// its 503 when the chi timeout fires; if it were shorter the
// server would abort the write first and the chi budget would
// be unreachable dead configuration.
writeTimeout = requestTimeout + 5*time.Second
)
func (s *Server) serveUntilShutdown() {
// newHTTPServer constructs the http.Server. It performs no I/O
// and does not start listening.
func (s *Server) newHTTPServer() *http.Server {
listenAddr := fmt.Sprintf(":%d", s.params.Config.Port)
s.httpServer = &http.Server{
return &http.Server{
Addr: listenAddr,
Handler: s,
MaxHeaderBytes: maxHeaderBytes,
ReadTimeout: readTimeout,
ReadHeaderTimeout: readHeaderTimeout,
WriteTimeout: writeTimeout,
IdleTimeout: idleTimeout,
}
}
s.SetupRoutes()
// listenAndServe runs the listener until the server is shut
// down. A genuine listen failure (not the expected
// ErrServerClosed from a clean shutdown) requests process
// shutdown through fx with a non-zero exit code, so the failure
// is visible to any supervisor.
func (s *Server) listenAndServe() {
s.log.Info("http begin listen",
"listenaddr", listenAddr,
"listenaddr", s.httpServer.Addr,
"version", s.params.Globals.Version,
"buildarch", s.params.Globals.Buildarch,
)
err := s.httpServer.ListenAndServe()
if err != nil && !errors.Is(err, http.ErrServerClosed) {
if err == nil || errors.Is(err, http.ErrServerClosed) {
return
}
s.log.Error("listen error", "error", err)
if s.cancelFunc != nil {
s.cancelFunc()
}
shutdownErr := s.shutdowner.Shutdown(fx.ExitCode(1))
if shutdownErr != nil {
s.log.Error("request shutdown failed", "error", shutdownErr)
}
}
+1
View File
@@ -17,6 +17,7 @@ func (s *Server) SetupRoutes() {
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(middleware.Timeout(requestTimeout))
+27 -71
View File
@@ -1,16 +1,14 @@
// Package server provides the HTTP server lifecycle,
// including startup, routing, signal handling, and graceful
// shutdown.
// including startup, routing, and graceful shutdown. The
// process lifetime is owned by fx: shutdown is requested
// through fx.Shutdowner so every component's OnStop hook runs
// in dependency order.
package server
import (
"context"
"log/slog"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"sneak.berlin/go/netwatch/internal/config"
"sneak.berlin/go/netwatch/internal/globals"
@@ -31,19 +29,18 @@ type Params struct {
Handlers *handlers.Handlers
Logger *logger.Logger
Middleware *middleware.Middleware
Shutdowner fx.Shutdowner
}
// Server is the top-level HTTP server orchestrator.
type Server struct {
cancelFunc context.CancelFunc
exitCode int
h *handlers.Handlers
httpServer *http.Server
log *slog.Logger
mw *middleware.Middleware
params Params
router *chi.Mux
startupTime time.Time
shutdowner fx.Shutdowner
}
// New creates a Server and registers lifecycle hooks for
@@ -57,23 +54,25 @@ func New(
s.mw = params.Middleware
s.h = params.Handlers
s.log = params.Logger.Get()
s.shutdowner = params.Shutdowner
lc.Append(fx.Hook{
OnStart: func(_ context.Context) error {
s.startupTime = time.Now().UTC()
// Build the router and http.Server synchronously
// here, before spawning the serving goroutine, so
// httpServer is fully constructed by the time OnStop
// (or an early signal) can read it. fx guarantees
// OnStart returns before OnStop runs, so no
// synchronization or nil check is needed at shutdown.
s.SetupRoutes()
s.httpServer = s.newHTTPServer()
go func() { //nolint:contextcheck // fx OnStart ctx is startup-only; run() creates its own
s.run()
}()
go s.listenAndServe()
return nil
},
OnStop: func(_ context.Context) error {
if s.cancelFunc != nil {
s.cancelFunc()
}
return nil
OnStop: func(ctx context.Context) error {
return s.shutdown(ctx)
},
})
@@ -88,60 +87,17 @@ func (s *Server) ServeHTTP(
s.router.ServeHTTP(w, r)
}
func (s *Server) run() {
exitCode := s.serve()
os.Exit(exitCode)
}
func (s *Server) serve() int {
var ctx context.Context //nolint:wsl // ctx must be declared before multi-assign
ctx, s.cancelFunc = context.WithCancel(
context.Background(),
)
go func() {
c := make(chan os.Signal, 1)
signal.Ignore(syscall.SIGPIPE)
signal.Notify(c, os.Interrupt, syscall.SIGTERM)
sig := <-c
s.log.Info("signal received", "signal", sig)
if s.cancelFunc != nil {
s.cancelFunc()
}
}()
go func() {
s.serveUntilShutdown()
}()
<-ctx.Done()
s.cleanShutdown()
return s.exitCode
}
const shutdownTimeout = 5 * time.Second
func (s *Server) cleanShutdown() {
s.exitCode = 0
ctxShutdown, shutdownCancel := context.WithTimeout(
context.Background(),
shutdownTimeout,
)
defer shutdownCancel()
err := s.httpServer.Shutdown(ctxShutdown)
// shutdown gracefully stops the HTTP server within the
// deadline of the context fx provides for OnStop.
func (s *Server) shutdown(ctx context.Context) error {
err := s.httpServer.Shutdown(ctx)
if err != nil {
s.log.Error(
"server clean shutdown failed",
"error", err,
)
s.log.Error("server clean shutdown failed", "error", err)
return err
}
s.log.Info("server stopped")
return nil
}