Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
fc511cd2d7 |
@@ -23,15 +23,6 @@ latest run passes.
|
|||||||
|
|
||||||
# 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
|
||||||
|
|||||||
@@ -40,12 +40,11 @@ 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
|
||||||
@@ -77,13 +76,8 @@ func New(
|
|||||||
return nil
|
return nil
|
||||||
},
|
},
|
||||||
OnStop: func(_ context.Context) error {
|
OnStop: func(_ context.Context) error {
|
||||||
// stopOnce makes OnStop idempotent: a second
|
close(b.done)
|
||||||
// invocation must not close an already-closed channel
|
b.flushLocked()
|
||||||
// (which would panic) or flush again.
|
|
||||||
b.stopOnce.Do(func() {
|
|
||||||
close(b.done)
|
|
||||||
b.flushLocked()
|
|
||||||
})
|
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -1,77 +1,13 @@
|
|||||||
package reportbuf_test
|
package reportbuf_test
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"os"
|
|
||||||
"strings"
|
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"sneak.berlin/go/netwatch/internal/config"
|
_ "sneak.berlin/go/netwatch/internal/reportbuf"
|
||||||
"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"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// TestFlushOnShutdown proves the flush-on-shutdown path: a
|
func TestImport(t *testing.T) {
|
||||||
// report appended after start but before the periodic flush
|
t.Parallel()
|
||||||
// window must reach disk when the fx lifecycle stops. This is
|
// Compilation check — verifies the package parses
|
||||||
// the exact case that silent data loss on restart used to
|
// and all imports resolve.
|
||||||
// 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,31 +5,20 @@ 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
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// newHTTPServer constructs the http.Server. It performs no I/O
|
func (s *Server) serveUntilShutdown() {
|
||||||
// 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)
|
||||||
|
|
||||||
return &http.Server{
|
s.httpServer = &http.Server{
|
||||||
Addr: listenAddr,
|
Addr: listenAddr,
|
||||||
Handler: s,
|
Handler: s,
|
||||||
MaxHeaderBytes: maxHeaderBytes,
|
MaxHeaderBytes: maxHeaderBytes,
|
||||||
@@ -38,29 +27,21 @@ func (s *Server) newHTTPServer() *http.Server {
|
|||||||
WriteTimeout: writeTimeout,
|
WriteTimeout: writeTimeout,
|
||||||
IdleTimeout: idleTimeout,
|
IdleTimeout: idleTimeout,
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
// listenAndServe runs the listener until the server is shut
|
s.SetupRoutes()
|
||||||
// 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", s.httpServer.Addr,
|
"listenaddr", listenAddr,
|
||||||
"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) {
|
||||||
return
|
s.log.Error("listen error", "error", err)
|
||||||
}
|
|
||||||
|
|
||||||
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,14 +1,16 @@
|
|||||||
// Package server provides the HTTP server lifecycle,
|
// Package server provides the HTTP server lifecycle,
|
||||||
// including startup, routing, and graceful shutdown. The
|
// including startup, routing, signal handling, and graceful
|
||||||
// process lifetime is owned by fx: shutdown is requested
|
// shutdown.
|
||||||
// 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"
|
||||||
@@ -29,18 +31,19 @@ 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 {
|
||||||
h *handlers.Handlers
|
cancelFunc context.CancelFunc
|
||||||
httpServer *http.Server
|
exitCode int
|
||||||
log *slog.Logger
|
h *handlers.Handlers
|
||||||
mw *middleware.Middleware
|
httpServer *http.Server
|
||||||
params Params
|
log *slog.Logger
|
||||||
router *chi.Mux
|
mw *middleware.Middleware
|
||||||
shutdowner fx.Shutdowner
|
params Params
|
||||||
|
router *chi.Mux
|
||||||
|
startupTime time.Time
|
||||||
}
|
}
|
||||||
|
|
||||||
// New creates a Server and registers lifecycle hooks for
|
// New creates a Server and registers lifecycle hooks for
|
||||||
@@ -54,25 +57,26 @@ 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 {
|
||||||
// Build the router and http.Server synchronously
|
s.startupTime = time.Now().UTC()
|
||||||
// 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 s.listenAndServe()
|
// The fx OnStart context is scoped to startup and is
|
||||||
|
// cancelled once the hook returns; run() derives its
|
||||||
|
// own context instead of inheriting this one.
|
||||||
|
go func() { //nolint:contextcheck // see comment above
|
||||||
|
s.run()
|
||||||
|
}()
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
},
|
},
|
||||||
OnStop: func(ctx context.Context) error {
|
OnStop: func(_ context.Context) error {
|
||||||
return s.shutdown(ctx)
|
if s.cancelFunc != nil {
|
||||||
|
s.cancelFunc()
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
|
|
||||||
@@ -87,17 +91,60 @@ func (s *Server) ServeHTTP(
|
|||||||
s.router.ServeHTTP(w, r)
|
s.router.ServeHTTP(w, r)
|
||||||
}
|
}
|
||||||
|
|
||||||
// shutdown gracefully stops the HTTP server within the
|
func (s *Server) run() {
|
||||||
// deadline of the context fx provides for OnStop.
|
exitCode := s.serve()
|
||||||
func (s *Server) shutdown(ctx context.Context) error {
|
os.Exit(exitCode)
|
||||||
err := s.httpServer.Shutdown(ctx)
|
}
|
||||||
if err != nil {
|
|
||||||
s.log.Error("server clean shutdown failed", "error", err)
|
|
||||||
|
|
||||||
return err
|
func (s *Server) serve() int {
|
||||||
|
var ctx context.Context
|
||||||
|
|
||||||
|
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 {
|
||||||
|
s.log.Error(
|
||||||
|
"server clean shutdown failed",
|
||||||
|
"error", err,
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
s.log.Info("server stopped")
|
s.log.Info("server stopped")
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user