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