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() + } + }) + } +}