Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
38 changes: 35 additions & 3 deletions server.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand Down
150 changes: 150 additions & 0 deletions start_failure_release_linux_test.go
Original file line number Diff line number Diff line change
@@ -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
}
Loading
Loading