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