diff --git a/TODO.md b/TODO.md index 29d1782..49f0492 100644 --- a/TODO.md +++ b/TODO.md @@ -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 diff --git a/backend/internal/reportbuf/reportbuf.go b/backend/internal/reportbuf/reportbuf.go index ff4cbad..1ef3b19 100644 --- a/backend/internal/reportbuf/reportbuf.go +++ b/backend/internal/reportbuf/reportbuf.go @@ -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 }, diff --git a/backend/internal/reportbuf/reportbuf_test.go b/backend/internal/reportbuf/reportbuf_test.go index 24a9ae9..4b6b550 100644 --- a/backend/internal/reportbuf/reportbuf_test.go +++ b/backend/internal/reportbuf/reportbuf_test.go @@ -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 } diff --git a/backend/internal/server/http.go b/backend/internal/server/http.go index a0d3c68..e621ce7 100644 --- a/backend/internal/server/http.go +++ b/backend/internal/server/http.go @@ -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) } } diff --git a/backend/internal/server/server.go b/backend/internal/server/server.go index a67bd1e..d14dd25 100644 --- a/backend/internal/server/server.go +++ b/backend/internal/server/server.go @@ -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 }