From 1a285c1334e7fa741990829790373a57cd6cc17f Mon Sep 17 00:00:00 2001 From: David Christopher <137860927+davidchrisr@users.noreply.github.com> Date: Thu, 17 Sep 2026 06:20:39 +0700 Subject: [PATCH] filer: join shutdown paths before closing metadata store (#11363) Serve can return when its listener closes while HTTP requests are still draining. The main path could then close the metadata store before those requests finish. Make signal, context, and Serve-exit paths join one shutdown sequence. Drain gRPC and HTTP concurrently with 15-second default limits, then close the store. Test both completion orders. --- weed/command/filer.go | 117 +++++++++----------- weed/command/filer_shutdown_test.go | 159 ++++++++++++++++++++++++++++ 2 files changed, 210 insertions(+), 66 deletions(-) create mode 100644 weed/command/filer_shutdown_test.go diff --git a/weed/command/filer.go b/weed/command/filer.go index 0fa629996..42099c19a 100644 --- a/weed/command/filer.go +++ b/weed/command/filer.go @@ -94,7 +94,7 @@ type FilerOptions struct { // and by weed mini; nil for standalone weed filer. shutdownCtx context.Context // gracefulStopTimeout caps how long startFiler waits for gRPC graceful - // stop before forcing the server to stop. Zero means the default of 10s. + // stop before forcing the server to stop. Zero means the default of 15s. gracefulStopTimeout time.Duration } @@ -416,14 +416,7 @@ func (fo *FilerOptions) startFiler() { // percent-encoded directory names when a client follows it (#11125). defaultHandler := weed_server.CleanPathHandler(defaultMux) - // Ensure fs.Shutdown() runs exactly once, whether triggered by a signal hook - // or by the main goroutine after Serve() returns (e.g., MiniCluster tests). - var shutdownOnce sync.Once - shutdownFiler := func() { - shutdownOnce.Do(func() { - fs.Shutdown() - }) - } + var httpServers []*http.Server if *fo.publicPort != 0 { publicListeningAddress := util.JoinHostPort(*fo.bindIp, *fo.publicPort) @@ -433,14 +426,18 @@ func (fo *FilerOptions) startFiler() { glog.Fatalf("Filer server public listener error on port %d:%v", *fo.publicPort, e) } publicHandler := weed_server.CleanPathHandler(publicVolumeMux) + publicServer := newHttpServer(publicHandler, nil) + httpServers = append(httpServers, publicServer) go func() { - if e := http.Serve(publicListener, publicHandler); e != nil { + if e := publicServer.Serve(publicListener); e != nil && e != http.ErrServerClosed { glog.Fatalf("Volume server fail to serve public: %v", e) } }() if localPublicListener != nil { + localPublicServer := newHttpServer(publicHandler, nil) + httpServers = append(httpServers, localPublicServer) go func() { - if e := http.Serve(localPublicListener, publicHandler); e != nil { + if e := localPublicServer.Serve(localPublicListener); e != nil && e != http.ErrServerClosed { glog.Errorf("Volume server fail to serve public: %v", e) } }() @@ -492,7 +489,7 @@ func (fo *FilerOptions) startFiler() { // Helper to gracefully stop the gRPC server, waiting for active RPCs. gracefulTimeout := fo.gracefulStopTimeout if gracefulTimeout <= 0 { - gracefulTimeout = 10 * time.Second + gracefulTimeout = 15 * time.Second } stopGrpcServer := func() { glog.V(0).Infof("Gracefully stopping gRPC server") @@ -524,6 +521,7 @@ func (fo *FilerOptions) startFiler() { glog.Fatalf("Failed to listen on %s: %v", localSocket, err) } socketServer = newHttpServer(defaultHandler, nil) + httpServers = append(httpServers, socketServer) go socketServer.Serve(filerSocketListener) } @@ -567,6 +565,7 @@ func (fo *FilerOptions) startFiler() { var localTLSServer *http.Server if filerLocalListener != nil { localTLSServer = newHttpServer(defaultHandler, tlsConfig) + httpServers = append(httpServers, localTLSServer) go func() { if err := localTLSServer.ServeTLS(filerLocalListener, "", ""); err != nil { glog.Errorf("Filer Fail to serve: %v", err) @@ -574,50 +573,28 @@ func (fo *FilerOptions) startFiler() { }() } httpS := newHttpServer(defaultHandler, tlsConfig) + httpServers = append(httpServers, httpS) + shutdown := newFilerShutdown(stopGrpcServer, fs.Shutdown, httpServers...) - // Register a single shutdown hook that runs the steps in the correct order: - // stop accepting new gRPC/HTTP requests, then close the filer database. - // Combining them into one hook keeps ordering intact regardless of how - // grace fires interrupt hooks (FIFO vs LIFO). - grace.OnInterrupt(func() { - stopGrpcServer() - glog.V(0).Infof("Gracefully stopping all HTTP servers") - shutdownCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second) - defer cancel() - if socketServer != nil { - err = socketServer.Shutdown(shutdownCtx) - if err != nil { - glog.Warningf("socket server shutdown: %v", err) - } - } - if localTLSServer != nil { - err = localTLSServer.Shutdown(shutdownCtx) - if err != nil { - glog.Warningf("local TLS server shutdown: %v", err) - } - } - if err := httpS.Shutdown(shutdownCtx); err != nil { - glog.Warningf("HTTPS server shutdown: %v", err) - } - shutdownFiler() - }) + grace.OnInterrupt(shutdown) if fo.shutdownCtx != nil { go func() { <-fo.shutdownCtx.Done() - httpS.Shutdown(context.Background()) - grpcS.Stop() + shutdown() }() } if err := httpS.ServeTLS(filerListener, "", ""); err != nil && err != http.ErrServerClosed { glog.Fatalf("Filer Fail to serve: %v", err) } - // Close database after servers have stopped to prevent data corruption - shutdownFiler() + // Serve returns when listeners close, before active requests finish. + // Join the same shutdown sequence instead of closing the filer here. + shutdown() } else { var localHTTPServer *http.Server if filerLocalListener != nil { localHTTPServer = newHttpServer(defaultHandler, nil) + httpServers = append(httpServers, localHTTPServer) go func() { if err := localHTTPServer.Serve(filerLocalListener); err != nil { glog.Errorf("Filer Fail to serve: %v", err) @@ -625,39 +602,47 @@ func (fo *FilerOptions) startFiler() { }() } httpS := newHttpServer(defaultHandler, nil) + httpServers = append(httpServers, httpS) + shutdown := newFilerShutdown(stopGrpcServer, fs.Shutdown, httpServers...) - // Register a single shutdown hook that runs the steps in the correct order: - // stop accepting new gRPC/HTTP requests, then close the filer database. - // Combining them into one hook keeps ordering intact regardless of how - // grace fires interrupt hooks (FIFO vs LIFO). - grace.OnInterrupt(func() { - stopGrpcServer() - glog.V(0).Infof("Gracefully stopping all HTTP servers") - shutdownCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second) - defer cancel() - if socketServer != nil { - socketServer.Shutdown(shutdownCtx) - } - if localHTTPServer != nil { - localHTTPServer.Shutdown(shutdownCtx) - } - if err := httpS.Shutdown(shutdownCtx); err != nil { - glog.Warningf("HTTP server shutdown: %v", err) - } - shutdownFiler() - }) + grace.OnInterrupt(shutdown) if fo.shutdownCtx != nil { go func() { <-fo.shutdownCtx.Done() - httpS.Shutdown(context.Background()) - grpcS.Stop() + shutdown() }() } if err := httpS.Serve(filerListener); err != nil && err != http.ErrServerClosed { glog.Fatalf("Filer Fail to serve: %v", err) } - // Close database after servers have stopped to prevent data corruption - shutdownFiler() + // Serve returns when listeners close, before active requests finish. + // Join the same shutdown sequence instead of closing the filer here. + shutdown() } } + +// newFilerShutdown joins shutdown callers while gRPC and HTTP drain concurrently. +func newFilerShutdown(stopGrpc, shutdownFiler func(), httpServers ...*http.Server) func() { + return sync.OnceFunc(func() { + shutdownCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + var drained sync.WaitGroup + drained.Add(1) + go func() { + defer drained.Done() + stopGrpc() + }() + for _, server := range httpServers { + drained.Add(1) + go func() { + defer drained.Done() + if err := server.Shutdown(shutdownCtx); err != nil { + glog.Warningf("filer HTTP shutdown: %v", err) + } + }() + } + drained.Wait() + shutdownFiler() + }) +} diff --git a/weed/command/filer_shutdown_test.go b/weed/command/filer_shutdown_test.go new file mode 100644 index 000000000..f115b8f61 --- /dev/null +++ b/weed/command/filer_shutdown_test.go @@ -0,0 +1,159 @@ +package command + +import ( + "io" + "net/http" + "net/http/httptest" + "sync" + "sync/atomic" + "testing" + "time" +) + +func TestFilerShutdownJoinsServeDuringParallelDrain(t *testing.T) { + var servers []*http.Server + var releases []func() + var completed []chan struct{} + var listenersClosing []chan struct{} + for range 2 { + entered := make(chan struct{}) + release := make(chan struct{}) + finish := sync.OnceFunc(func() { close(release) }) + done := make(chan struct{}) + closing := make(chan struct{}) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + close(entered) + <-release + w.WriteHeader(http.StatusNoContent) + })) + t.Cleanup(server.Close) + t.Cleanup(finish) + server.Config.RegisterOnShutdown(func() { close(closing) }) + servers = append(servers, server.Config) + releases = append(releases, finish) + completed = append(completed, done) + listenersClosing = append(listenersClosing, closing) + go func() { + defer close(done) + response, err := server.Client().Get(server.URL) + if err != nil { + t.Error(err) + return + } + defer response.Body.Close() + io.Copy(io.Discard, response.Body) + if response.StatusCode != http.StatusNoContent { + t.Errorf("active request returned %s", response.Status) + } + }() + <-entered + } + + grpcStarted := make(chan struct{}) + grpcRelease := make(chan struct{}) + releaseGrpc := sync.OnceFunc(func() { close(grpcRelease) }) + t.Cleanup(releaseGrpc) + grpcStopped := make(chan struct{}) + filerClosed := make(chan struct{}) + var closes atomic.Int32 + shutdown := newFilerShutdown(func() { + close(grpcStarted) + <-grpcRelease + close(grpcStopped) + }, func() { + closes.Add(1) + close(filerClosed) + }, servers...) + joined := make(chan struct{}) + go func() { shutdown(); close(joined) }() + deadline := time.After(5 * time.Second) + select { + case <-grpcStarted: + case <-deadline: + t.Fatal("gRPC shutdown did not start") + } + for _, closing := range listenersClosing { + select { + case <-closing: + case <-deadline: + t.Fatal("HTTP shutdown did not start while gRPC was draining") + } + } + // Model the main Serve path joining shutdown as soon as listeners close. + serveReturned := make(chan struct{}) + go func() { shutdown(); close(serveReturned) }() + + releases[0]() + <-completed[0] + select { + case <-filerClosed: + t.Error("filer closed while another server still had an active request") + default: + } + select { + case <-serveReturned: + t.Error("Serve exit path returned before request draining completed") + default: + } + releases[1]() + <-completed[1] + select { + case <-filerClosed: + t.Error("filer closed while gRPC was still draining") + default: + } + releaseGrpc() + <-grpcStopped + <-joined + <-serveReturned + if closes.Load() != 1 { + t.Errorf("filer closed %d times, want once", closes.Load()) + } +} + +func TestFilerShutdownWaitsForHTTPAfterGrpcStops(t *testing.T) { + entered := make(chan struct{}) + release := make(chan struct{}) + finish := sync.OnceFunc(func() { close(release) }) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + close(entered) + <-release + w.WriteHeader(http.StatusNoContent) + })) + t.Cleanup(server.Close) + t.Cleanup(finish) + requestDone := make(chan struct{}) + go func() { + defer close(requestDone) + response, err := server.Client().Get(server.URL) + if err != nil { + t.Error(err) + return + } + response.Body.Close() + }() + <-entered + + httpClosing := make(chan struct{}) + server.Config.RegisterOnShutdown(func() { close(httpClosing) }) + grpcStopped := make(chan struct{}) + filerClosed := make(chan struct{}) + shutdown := newFilerShutdown(func() { close(grpcStopped) }, func() { close(filerClosed) }, server.Config) + joined := make(chan struct{}) + go func() { shutdown(); close(joined) }() + <-httpClosing + <-grpcStopped + select { + case <-filerClosed: + t.Error("filer closed while HTTP request was active") + case <-time.After(100 * time.Millisecond): + } + finish() + <-requestDone + <-joined + select { + case <-filerClosed: + default: + t.Error("filer was not closed after HTTP request completed") + } +}