main into prod: netwatch as one container, ready for upaas #74

Open
clawbot wants to merge 25 commits from main into prod
5 changed files with 154 additions and 100 deletions
Showing only changes of commit d7cf010e00 - Show all commits
+9
View File
@@ -22,6 +22,15 @@ 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
+13 -7
View File
@@ -40,11 +40,12 @@ type Params struct {
// Buffer accumulates JSON lines in memory and flushes them
// to zstd-compressed files on disk.
type Buffer struct {
buf bytes.Buffer
dataDir string
done chan struct{}
log *slog.Logger
mu sync.Mutex
buf bytes.Buffer
dataDir string
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 {
close(b.done)
b.flushLocked()
// 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
}
+30 -11
View File
@@ -5,20 +5,31 @@ import (
"fmt"
"net/http"
"time"
"go.uber.org/fx"
)
const (
readTimeout = 10 * time.Second
readHeaderTimeout = 5 * time.Second
writeTimeout = 10 * 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,
@@ -27,21 +38,29 @@ func (s *Server) serveUntilShutdown() {
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) {
s.log.Error("listen error", "error", err)
if err == nil || errors.Is(err, http.ErrServerClosed) {
return
}
if s.cancelFunc != nil {
s.cancelFunc()
}
s.log.Error("listen error", "error", err)
shutdownErr := s.shutdowner.Shutdown(fx.ExitCode(1))
if shutdownErr != nil {
s.log.Error("request shutdown failed", "error", shutdownErr)
}
}
+33 -77
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
h *handlers.Handlers
httpServer *http.Server
log *slog.Logger
mw *middleware.Middleware
params Params
router *chi.Mux
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
}