Compare commits
1
Commits
924d4b7177
...
dd762ce759
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
dd762ce759 |
@@ -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
|
||||
|
||||
@@ -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
|
||||
},
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
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,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
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user