From fb52752977153881b94992b18a7d987aba1d88c4 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 20:12:35 +0200 Subject: [PATCH] fix(server): a Start that never serves releases the caller's listener, the CPU monitor and the settle re-opener (celeris#737) A Start that failed in doPrepare returned with the listener the caller had handed over still bound and listening, and, when createEngine failed, with the CPU monitor's /proc/stat descriptor open: only Shutdown closed it, and such a caller has no reason to call Shutdown. A Start that published its engine but ended with no Shutdown to come, because Listen failed or Shutdown had been called before Start, likewise left the monitor open and the settle re-opener (celeris#592) running. prepareWithListener now closes the supplied listener on any start failure except ErrAlreadyStarted, doPrepare closes the monitor when createEngine fails, and the function listenContext returns, which every Start* entry point defers until Listen has returned, also stops the re-opener and closes the monitor. All three releases are idempotent. --- server.go | 38 +++- start_failure_release_linux_test.go | 150 ++++++++++++++++ start_failure_release_test.go | 260 ++++++++++++++++++++++++++++ 3 files changed, 445 insertions(+), 3 deletions(-) create mode 100644 start_failure_release_linux_test.go create mode 100644 start_failure_release_test.go diff --git a/server.go b/server.go index 8985b94b..72ff89ac 100644 --- a/server.go +++ b/server.go @@ -403,6 +403,16 @@ func (s *Server) Start() error { // cancels Listen exactly as before, this only adds a second way to wake it. // If Shutdown already ran (or is racing prepare), the returned context is // already cancelled so Listen returns immediately instead of parking forever. +// +// The returned function cancels that context and releases what doPrepare +// opened for the run: the settle re-opener and the CPU monitor. Every Start* +// entry point defers it right after this call, so it runs once Listen has +// returned, and from then on the server never serves again (Start cannot be +// retried). Shutdown releases the same two, but a Start that ends without a +// Shutdown to come, because Listen failed or because Shutdown was called +// before Start, left the re-opener running and the monitor's /proc/stat +// descriptor open for the life of the process (celeris#737). Both releases +// are idempotent, so a Shutdown before or after it is harmless. func (s *Server) listenContext(parent context.Context) (context.Context, context.CancelFunc) { ctx, cancel := context.WithCancel(parent) s.lifecycleMu.Lock() @@ -414,7 +424,11 @@ func (s *Server) listenContext(parent context.Context) (context.Context, context if alreadyShutdown { cancel() } - return ctx, cancel + return ctx, func() { + cancel() + s.router.stopSettleReopener() + s.closeCPUMonitor() + } } // cancelListen wakes a Listen parked on the context published by @@ -684,9 +698,19 @@ func (s *Server) logger() *slog.Logger { } func (s *Server) prepareWithListener(ln net.Listener) (engine.Engine, error) { - return s.doPrepare(func(cfg *resource.Config) { + eng, err := s.doPrepare(func(cfg *resource.Config) { cfg.Listener = ln }) + // celeris#737: the caller handed ln over and may not close it, so a start + // that fails before any engine runs closes it here. Otherwise it stays + // bound, and the kernel keeps completing handshakes into a backlog that + // nothing accepts. ErrAlreadyStarted is the exception: a server is + // already running, and ln may be the very listener it serves on, so it + // stays the caller's. + if err != nil && !errors.Is(err, ErrAlreadyStarted) && ln != nil { + _ = ln.Close() + } + return eng, err } // doPrepare is the shared implementation for prepare and prepareWithListener. @@ -781,6 +805,10 @@ func (s *Server) doPrepare(configureFn func(cfg *resource.Config)) (engine.Engin eng, err = createEngine(cfg, handler, cpuMon) if err != nil { s.startErr = fmt.Errorf("create engine: %w", err) + // celeris#737: only Shutdown closes the monitor, and a caller + // whose Start failed has no reason to call it, so the + // /proc/stat descriptor opened above would stay open. + s.closeCPUMonitor() return } s.engineRef.Store(&eng) @@ -843,7 +871,11 @@ func (s *Server) doPrepare(configureFn func(cfg *resource.Config)) (engine.Engin // extracted from ln and the listener is closed so the engine workers can // rebind their own SO_REUSEPORT sockets bound to the same (host, port). // In both cases, the caller must not Accept on or close the supplied -// listener after calling this function. +// listener after calling this function. If the server fails to start before +// its engine runs (a configuration error or an engine that cannot be +// created), the listener is closed before the error is returned. The one +// exception is [ErrAlreadyStarted]: the server is already running, perhaps on +// that very listener, so this call leaves it to the caller. // // On the adaptive engine (the default on Linux) the listener goes to whichever // sub-engine starts, and a later switch binds the second sub-engine to that diff --git a/start_failure_release_linux_test.go b/start_failure_release_linux_test.go new file mode 100644 index 00000000..a3cd3e3d --- /dev/null +++ b/start_failure_release_linux_test.go @@ -0,0 +1,150 @@ +//go:build linux + +package celeris_test + +import ( + "context" + "io" + "log/slog" + "net" + "os" + "path/filepath" + "runtime/debug" + "testing" + + "github.com/goceleris/celeris" +) + +// TestFailedStartsLeaveNoProcStatDescriptor is celeris#737 as it was +// measured: twenty starts that fail to create their engine, with GC off so no +// finalizer can close a leaked descriptor, must leave no /proc/stat +// descriptor open. doPrepare opens one for the CPU monitor before it creates +// the engine, and only Shutdown closed it. The count is of descriptors whose +// link is /proc/stat, so descriptors other tests open or close meanwhile +// cannot move it. +func TestFailedStartsLeaveNoProcStatDescriptor(t *testing.T) { + const starts = 20 + prev := debug.SetGCPercent(-1) + defer debug.SetGCPercent(prev) + quiet := slog.New(slog.NewTextHandler(io.Discard, nil)) + + entries := []struct { + name string + start func(s *celeris.Server) error + }{ + {"Start", func(s *celeris.Server) error { return s.Start() }}, + {"StartWithContext", func(s *celeris.Server) error { return s.StartWithContext(context.Background()) }}, + {"StartWithListener", func(s *celeris.Server) error { + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + return err + } + return s.StartWithListener(ln) + }}, + {"StartWithListenerAndContext", func(s *celeris.Server) error { + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + return err + } + return s.StartWithListenerAndContext(context.Background(), ln) + }}, + } + for _, entry := range entries { + t.Run(entry.name, func(t *testing.T) { + before := procStatDescriptors(t) + for range starts { + // No factory knows engine type 99, so createEngine fails. + s := celeris.New(celeris.Config{Engine: celeris.EngineType(99), Addr: "127.0.0.1:0", Logger: quiet}) + if err := entry.start(s); err == nil { + t.Fatalf("%s with an unknown engine type returned nil", entry.name) + } + } + after := procStatDescriptors(t) + t.Logf("%s: /proc/stat descriptors %d before %d failed starts, %d after", entry.name, before, starts, after) + if after != before { + t.Errorf("%s: %d failed starts left %d /proc/stat descriptors open (celeris#737)", entry.name, starts, after-before) + } + }) + } +} + +// TestStartsWithoutShutdownLeaveNoProcStatDescriptor is the rest of the +// celeris#737 class, measured the same way: Starts that publish their engine +// and run Listen but end with no Shutdown to come, because Listen fails (the +// address is taken by a socket without SO_REUSEPORT) or because Shutdown was +// called before Start. On std and on epoll, which binds per worker. +func TestStartsWithoutShutdownLeaveNoProcStatDescriptor(t *testing.T) { + const starts = 20 + prev := debug.SetGCPercent(-1) + defer debug.SetGCPercent(prev) + quiet := slog.New(slog.NewTextHandler(io.Discard, nil)) + + busy, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + defer func() { _ = busy.Close() }() + busyAddr := busy.Addr().String() + + cases := []struct { + name string + cfg celeris.Config + shutdownFirst bool + wantErr bool + start func(s *celeris.Server) error + }{ + {"Listen fails/std/Start", celeris.Config{Engine: celeris.Std, Addr: busyAddr}, false, true, + func(s *celeris.Server) error { return s.Start() }}, + {"Listen fails/epoll/StartWithContext", celeris.Config{Engine: celeris.Epoll, Addr: busyAddr, Workers: 2}, false, true, + func(s *celeris.Server) error { return s.StartWithContext(context.Background()) }}, + {"Shutdown first/std/StartWithContext", celeris.Config{Engine: celeris.Std, Addr: "127.0.0.1:0"}, true, false, + func(s *celeris.Server) error { return s.StartWithContext(context.Background()) }}, + {"Shutdown first/epoll/StartWithListener", celeris.Config{Engine: celeris.Epoll, Workers: 2}, true, false, + func(s *celeris.Server) error { + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + return err + } + return s.StartWithListener(ln) + }}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + before := procStatDescriptors(t) + for i := range starts { + cfg := tc.cfg + cfg.Logger = quiet + s := celeris.New(cfg) + if tc.shutdownFirst { + if err := s.Shutdown(context.Background()); err != nil { + t.Fatalf("Shutdown before Start: %v", err) + } + } + if err := tc.start(s); (err != nil) != tc.wantErr { + t.Fatalf("start %d: Start returned %v, want an error: %v", i, err, tc.wantErr) + } + } + after := procStatDescriptors(t) + t.Logf("%s: /proc/stat descriptors %d before %d Starts, %d after", tc.name, before, starts, after) + if after != before { + t.Errorf("%s: %d Starts with no Shutdown to come left %d /proc/stat descriptors open (celeris#737)", tc.name, starts, after-before) + } + }) + } +} + +// procStatDescriptors counts this process's descriptors open on /proc/stat. +func procStatDescriptors(t *testing.T) int { + t.Helper() + ents, err := os.ReadDir("/proc/self/fd") + if err != nil { + t.Fatalf("read /proc/self/fd: %v", err) + } + n := 0 + for _, e := range ents { + if target, err := os.Readlink(filepath.Join("/proc/self/fd", e.Name())); err == nil && target == "/proc/stat" { + n++ + } + } + return n +} diff --git a/start_failure_release_test.go b/start_failure_release_test.go new file mode 100644 index 00000000..976f4227 --- /dev/null +++ b/start_failure_release_test.go @@ -0,0 +1,260 @@ +package celeris + +import ( + "context" + "errors" + "io" + "log/slog" + "net" + "testing" + "time" + + "github.com/goceleris/celeris/engine" +) + +// startFailures are the ways a Start can fail before any engine runs: in +// doPrepare, which returns before Listen is ever called. +func startFailures() []struct { + name string + cfg Config +} { + return []struct { + name string + cfg Config + }{ + {"config validation", Config{Engine: Std, MaxHeaderBytes: 100}}, + {"TrustedProxies", Config{Engine: Std, TrustedProxies: []string{"not-an-ip"}}}, + // No factory knows this engine type, so createEngine fails, after + // doPrepare has opened the CPU monitor. + {"createEngine", Config{Engine: EngineType(99)}}, + } +} + +// listenerStarts are the entry points that take a caller's listener. +func listenerStarts() []struct { + name string + start func(s *Server, ln net.Listener) error +} { + return []struct { + name string + start func(s *Server, ln net.Listener) error + }{ + {"StartWithListener", func(s *Server, ln net.Listener) error { return s.StartWithListener(ln) }}, + {"StartWithListenerAndContext", func(s *Server, ln net.Listener) error { + return s.StartWithListenerAndContext(context.Background(), ln) + }}, + } +} + +// TestFailedStartClosesTheSuppliedListener pins celeris#737. StartWithListener +// and StartWithListenerAndContext tell the caller not to close the listener +// it hands over, and a start that failed before its engine ran returned the +// error with that listener still open: bound and listening, so the kernel kept +// completing handshakes into a backlog nothing would ever accept. It must be +// closed, and a client dialling its address refused. +// +// A second call on the same server returns the first call's error without +// running doPrepare again, and must close the listener it is handed as well. +func TestFailedStartClosesTheSuppliedListener(t *testing.T) { + quiet := slog.New(slog.NewTextHandler(io.Discard, nil)) + for _, f := range startFailures() { + for _, entry := range listenerStarts() { + t.Run(f.name+"/"+entry.name, func(t *testing.T) { + cfg := f.cfg + cfg.Logger = quiet + s := New(cfg) + for call := 1; call <= 2; call++ { + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + addr := ln.Addr().String() + startErr := entry.start(s, ln) + if startErr == nil || errors.Is(startErr, ErrAlreadyStarted) { + _ = ln.Close() + t.Fatalf("call %d: %s returned %v, want the start failure", call, entry.name, startErr) + } + // Closing a listener that is already closed reports + // net.ErrClosed; closing an open one succeeds (and tidies up). + if cerr := ln.Close(); !errors.Is(cerr, net.ErrClosed) { + t.Errorf("call %d: %s failed (%v) and left the supplied listener open (celeris#737)", call, entry.name, startErr) + } + if c, derr := net.DialTimeout("tcp", addr, time.Second); derr == nil { + _ = c.Close() + t.Errorf("call %d: after the failed start a dial to the listener's address %s still connected", call, addr) + } + } + }) + } + } +} + +// TestErrAlreadyStartedLeavesTheListenerToTheCaller is the boundary of the +// celeris#737 fix: a call that finds the server already started does not own +// the listener it was handed. It may be the listener the running server +// serves on (std uses it directly), so closing it would stop that server. +func TestErrAlreadyStartedLeavesTheListenerToTheCaller(t *testing.T) { + for _, entry := range listenerStarts() { + t.Run(entry.name, func(t *testing.T) { + s := New(Config{}) + s.startOnce.Do(func() { + var fe engine.Engine = &fakeEngine{} + s.engineRef.Store(&fe) + }) + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + if err := entry.start(s, ln); !errors.Is(err, ErrAlreadyStarted) { + _ = ln.Close() + t.Fatalf("%s on a started server returned %v, want ErrAlreadyStarted", entry.name, err) + } + if cerr := ln.Close(); cerr != nil { + t.Errorf("%s returned ErrAlreadyStarted and closed the caller's listener (Close: %v)", entry.name, cerr) + } + }) + } +} + +// TestFailedStartClosesTheCPUMonitor pins the other half of celeris#737: +// doPrepare opens the CPU monitor (a /proc/stat descriptor on Linux) before it +// creates the engine, and only Shutdown closed it, which a caller whose Start +// failed has no reason to call. Every entry point shares doPrepare. +func TestFailedStartClosesTheCPUMonitor(t *testing.T) { + quiet := slog.New(slog.NewTextHandler(io.Discard, nil)) + entries := []struct { + name string + start func(s *Server) error + }{ + {"Start", func(s *Server) error { return s.Start() }}, + {"StartWithContext", func(s *Server) error { return s.StartWithContext(context.Background()) }}, + {"StartWithListener", func(s *Server) error { + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + return err + } + return s.StartWithListener(ln) + }}, + } + for _, entry := range entries { + t.Run(entry.name, func(t *testing.T) { + s := New(Config{Engine: EngineType(99), Addr: "127.0.0.1:0", Logger: quiet}) + err := entry.start(s) + if err == nil { + t.Fatalf("%s with an unknown engine type returned nil", entry.name) + } + s.cpuMonMu.Lock() + mon := s.cpuMon + s.cpuMonMu.Unlock() + if mon != nil { + t.Errorf("%s failed to create its engine (%v) and left the CPU monitor open (celeris#737)", entry.name, err) + s.closeCPUMonitor() + } + }) + } +} + +// TestStartWithoutShutdownReleasesWhatItOpened covers the rest of the +// celeris#737 class: a Start that published its engine and ran Listen, but +// ends with no Shutdown to come. Listen failed (the address is taken), or +// Shutdown was called before Start, so Listen ran on a context that was +// already cancelled and returned at once. doPrepare had opened the CPU monitor +// and started the settle re-opener (celeris#592), and only Shutdown released +// them: a descriptor and a goroutine per such Start, for the life of the +// process. Every Start* entry point now releases both once Listen has +// returned. +func TestStartWithoutShutdownReleasesWhatItOpened(t *testing.T) { + quiet := slog.New(slog.NewTextHandler(io.Discard, nil)) + busy, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + defer func() { _ = busy.Close() }() + busyAddr := busy.Addr().String() + + closedListener := func(t *testing.T) net.Listener { + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + _ = ln.Close() + return ln + } + openListener := func(t *testing.T) net.Listener { + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + return ln + } + cases := []struct { + name string + addr string + // listener, when set, makes the listener handed to start. + listener func(t *testing.T) net.Listener + // shutdownFirst: Shutdown is called before Start. + shutdownFirst bool + start func(s *Server, ln net.Listener) error + wantErr bool + }{ + {"Listen fails/Start", busyAddr, nil, false, func(s *Server, _ net.Listener) error { return s.Start() }, true}, + {"Listen fails/StartWithContext", busyAddr, nil, false, func(s *Server, _ net.Listener) error { + return s.StartWithContext(context.Background()) + }, true}, + {"Listen fails/StartWithListener", "", closedListener, false, func(s *Server, ln net.Listener) error { + return s.StartWithListener(ln) + }, true}, + {"Shutdown first/Start", "127.0.0.1:0", nil, true, func(s *Server, _ net.Listener) error { return s.Start() }, false}, + {"Shutdown first/StartWithContext", "127.0.0.1:0", nil, true, func(s *Server, _ net.Listener) error { + return s.StartWithContext(context.Background()) + }, false}, + {"Shutdown first/StartWithListenerAndContext", "", openListener, true, func(s *Server, ln net.Listener) error { + return s.StartWithListenerAndContext(context.Background(), ln) + }, false}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + // AsyncHandlers makes /s an adaptive route, and only adaptive + // routes start the settle re-opener. + s := New(Config{Engine: Std, Addr: tc.addr, AsyncHandlers: true, Logger: quiet}) + s.GET("/s", noopHandler) + if !s.router.adaptiveRoutes["/s"] { + t.Fatal("precondition: /s must be adaptive, or the re-opener never starts") + } + if tc.shutdownFirst { + if err := s.Shutdown(context.Background()); err != nil { + t.Fatalf("Shutdown before Start: %v", err) + } + } + var ln net.Listener + if tc.listener != nil { + ln = tc.listener(t) + } + done := make(chan error, 1) + go func() { done <- tc.start(s, ln) }() + var err error + select { + case err = <-done: + case <-time.After(10 * time.Second): + t.Fatal("Start did not return") + } + if tc.wantErr != (err != nil) { + t.Fatalf("Start returned %v, want an error: %v", err, tc.wantErr) + } + if s.loadEngine() == nil { + t.Fatal("precondition: the engine was never published, so this is not the case under test") + } + if reopenerRunning(s.router) { + t.Errorf("Start returned (%v) with the settle re-opener still running and no Shutdown to come (celeris#737)", err) + s.router.stopSettleReopener() + } + s.cpuMonMu.Lock() + mon := s.cpuMon + s.cpuMonMu.Unlock() + if mon != nil { + t.Errorf("Start returned (%v) with the CPU monitor still open and no Shutdown to come (celeris#737)", err) + s.closeCPUMonitor() + } + }) + } +}