Author SHA1 Message Date
sneak dd762ce759 fix(server): shut down through fx so buffered reports flush (closes #22)
check / check (push) Failing after 0s
The server ran os.Exit at the end of its own goroutine, which raced fx's
teardown and could kill the process before reportbuf's OnStop flushed the
buffer — losing up to a full flush window of telemetry on every restart,
silently and with exit 0. The server now requests shutdown through
fx.Shutdowner, so fx runs every OnStop in dependency order.

The http.Server is built synchronously in OnStart before the serving
goroutine starts, so shutdown can no longer race or nil-deref it; the
field is never written and read from two goroutines without a
happens-before edge. A listen failure now exits non-zero via
fx.ExitCode(1). reportbuf's OnStop is guarded by sync.Once. writeTimeout
now exceeds the chi per-request budget, with a comment, so that budget is
reachable. Dead startupTime, exitCode, and cancelFunc fields are gone.

A new test buffers a report and asserts it reaches disk after the fx
lifecycle stops.

Model: opus-4-8
2026-09-21 12:49:59 +00:00
5 changed files with 154 additions and 100 deletions
+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-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
+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,39 +5,58 @@ import (
"fmt"
"net/http"
"time"
"go.uber.org/fx"
)
const (
readTimeout = 10 * time.Second
writeTimeout = 10 * 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,
WriteTimeout: writeTimeout,
}
}
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
}