fix(server): shut down through fx so buffered reports flush #57

Merged
clawbot merged 1 commits from fix/shutdown-lifecycle into next 2026-09-21 18:47:12 +02:00
5 changed files with 154 additions and 100 deletions
Showing only changes of commit 924d4b7177 - Show all commits
+9
View File
@@ -22,6 +22,15 @@ files, so merging it also closes most compliance gaps.
# Completed Steps # 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 - 2026-09-21: backend HTTP hardening (issue #19): added `ReadHeaderTimeout` and
`IdleTimeout` to the server, a `SecurityHeaders` middleware (HSTS, tight CSP, `IdleTimeout` to the server, a `SecurityHeaders` middleware (HSTS, tight CSP,
frame/sniff/referrer/permissions headers) registered before CORS, and 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 // Buffer accumulates JSON lines in memory and flushes them
// to zstd-compressed files on disk. // to zstd-compressed files on disk.
type Buffer struct { type Buffer struct {
buf bytes.Buffer buf bytes.Buffer
dataDir string dataDir string
done chan struct{} done chan struct{}
log *slog.Logger log *slog.Logger
mu sync.Mutex mu sync.Mutex
stopOnce sync.Once
} }
// New creates a Buffer and registers lifecycle hooks to // New creates a Buffer and registers lifecycle hooks to
@@ -76,8 +77,13 @@ func New(
return nil return nil
}, },
OnStop: func(_ context.Context) error { OnStop: func(_ context.Context) error {
close(b.done) // stopOnce makes OnStop idempotent: a second
b.flushLocked() // 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 return nil
}, },
+69 -5
View File
@@ -1,13 +1,77 @@
package reportbuf_test package reportbuf_test
import ( import (
"os"
"strings"
"testing" "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) { // TestFlushOnShutdown proves the flush-on-shutdown path: a
t.Parallel() // report appended after start but before the periodic flush
// Compilation check — verifies the package parses // window must reach disk when the fx lifecycle stops. This is
// and all imports resolve. // 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" "fmt"
"net/http" "net/http"
"time" "time"
"go.uber.org/fx"
) )
const ( const (
readTimeout = 10 * time.Second readTimeout = 10 * time.Second
readHeaderTimeout = 5 * time.Second readHeaderTimeout = 5 * time.Second
writeTimeout = 10 * time.Second
idleTimeout = 60 * time.Second idleTimeout = 60 * time.Second
maxHeaderBytes = 1 << 20 // 1 MiB 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) listenAddr := fmt.Sprintf(":%d", s.params.Config.Port)
s.httpServer = &http.Server{ return &http.Server{
Addr: listenAddr, Addr: listenAddr,
Handler: s, Handler: s,
MaxHeaderBytes: maxHeaderBytes, MaxHeaderBytes: maxHeaderBytes,
@@ -27,21 +38,29 @@ func (s *Server) serveUntilShutdown() {
WriteTimeout: writeTimeout, WriteTimeout: writeTimeout,
IdleTimeout: idleTimeout, 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", s.log.Info("http begin listen",
"listenaddr", listenAddr, "listenaddr", s.httpServer.Addr,
"version", s.params.Globals.Version, "version", s.params.Globals.Version,
"buildarch", s.params.Globals.Buildarch, "buildarch", s.params.Globals.Buildarch,
) )
err := s.httpServer.ListenAndServe() err := s.httpServer.ListenAndServe()
if err != nil && !errors.Is(err, http.ErrServerClosed) { if err == nil || errors.Is(err, http.ErrServerClosed) {
s.log.Error("listen error", "error", err) return
}
if s.cancelFunc != nil { s.log.Error("listen error", "error", err)
s.cancelFunc()
} 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, // Package server provides the HTTP server lifecycle,
// including startup, routing, signal handling, and graceful // including startup, routing, and graceful shutdown. The
// shutdown. // process lifetime is owned by fx: shutdown is requested
// through fx.Shutdowner so every component's OnStop hook runs
// in dependency order.
package server package server
import ( import (
"context" "context"
"log/slog" "log/slog"
"net/http" "net/http"
"os"
"os/signal"
"syscall"
"time"
"sneak.berlin/go/netwatch/internal/config" "sneak.berlin/go/netwatch/internal/config"
"sneak.berlin/go/netwatch/internal/globals" "sneak.berlin/go/netwatch/internal/globals"
@@ -31,19 +29,18 @@ type Params struct {
Handlers *handlers.Handlers Handlers *handlers.Handlers
Logger *logger.Logger Logger *logger.Logger
Middleware *middleware.Middleware Middleware *middleware.Middleware
Shutdowner fx.Shutdowner
} }
// Server is the top-level HTTP server orchestrator. // Server is the top-level HTTP server orchestrator.
type Server struct { type Server struct {
cancelFunc context.CancelFunc h *handlers.Handlers
exitCode int httpServer *http.Server
h *handlers.Handlers log *slog.Logger
httpServer *http.Server mw *middleware.Middleware
log *slog.Logger params Params
mw *middleware.Middleware router *chi.Mux
params Params shutdowner fx.Shutdowner
router *chi.Mux
startupTime time.Time
} }
// New creates a Server and registers lifecycle hooks for // New creates a Server and registers lifecycle hooks for
@@ -57,23 +54,25 @@ func New(
s.mw = params.Middleware s.mw = params.Middleware
s.h = params.Handlers s.h = params.Handlers
s.log = params.Logger.Get() s.log = params.Logger.Get()
s.shutdowner = params.Shutdowner
lc.Append(fx.Hook{ lc.Append(fx.Hook{
OnStart: func(_ context.Context) error { 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 go s.listenAndServe()
s.run()
}()
return nil return nil
}, },
OnStop: func(_ context.Context) error { OnStop: func(ctx context.Context) error {
if s.cancelFunc != nil { return s.shutdown(ctx)
s.cancelFunc()
}
return nil
}, },
}) })
@@ -88,60 +87,17 @@ func (s *Server) ServeHTTP(
s.router.ServeHTTP(w, r) s.router.ServeHTTP(w, r)
} }
func (s *Server) run() { // shutdown gracefully stops the HTTP server within the
exitCode := s.serve() // deadline of the context fx provides for OnStop.
os.Exit(exitCode) func (s *Server) shutdown(ctx context.Context) error {
} err := s.httpServer.Shutdown(ctx)
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)
if err != nil { if err != nil {
s.log.Error( s.log.Error("server clean shutdown failed", "error", err)
"server clean shutdown failed",
"error", err, return err
)
} }
s.log.Info("server stopped") s.log.Info("server stopped")
return nil
} }