MCPcopy Create free account
hub / github.com/PostHog/duckgres / ShutdownAll

Method ShutdownAll

controlplane/worker_mgr.go:1045–1096  ·  view source on GitHub ↗

ShutdownAll stops all workers gracefully.

()

Source from the content-addressed store, hash-verified

1043 case <-time.After(3 * time.Second):
1044 slog.Warn("Worker did not exit in time, killing.", "id", w.ID)
1045 if w.cmd != nil && w.cmd.Process != nil {
1046 _ = w.cmd.Process.Kill()
1047 }
1048 if w.done != nil {
1049 <-w.done
1050 }
1051 }
1052 }
1053
1054 // Close gRPC client after the process has exited
1055 if w.client != nil {
1056 _ = w.client.Close()
1057 }
1058
1059 // Return pre-bound socket to pool for reuse, or close non-pre-bound listener.
1060 p.releaseWorkerSocket(w)
1061}
1062
1063// SetMaxWorkers updates the maximum number of workers. 0 means unlimited.
1064func (p *FlightWorkerPool) SetMaxWorkers(n int) {
1065 p.mu.Lock()
1066 defer p.mu.Unlock()
1067 p.maxWorkers = n
1068}
1069
1070// ShutdownAll stops all workers gracefully.
1071func (p *FlightWorkerPool) ShutdownAll() {
1072 p.mu.Lock()
1073 if p.shuttingDown {
1074 p.mu.Unlock()
1075 return
1076 }
1077 p.shuttingDown = true
1078 workers := make([]*ManagedWorker, 0, len(p.workers))
1079 for _, w := range p.workers {
1080 workers = append(workers, w)
1081 }
1082 p.mu.Unlock()
1083
1084 close(p.shutdownCh)
1085
1086 for _, w := range workers {
1087 if w.cmd.Process != nil {
1088 slog.Info("Shutting down worker.", "id", w.ID, "pid", w.cmd.Process.Pid)
1089 _ = w.cmd.Process.Signal(os.Interrupt)
1090 }
1091 }
1092
1093 // Wait up to 10s for workers to exit
1094 for _, w := range workers {
1095 select {
1096 case <-w.done:
1097 case <-time.After(10 * time.Second):
1098 slog.Warn("Worker did not exit in time, killing.", "id", w.ID)
1099 if w.cmd.Process != nil {

Calls 4

releaseWorkerSocketMethod · 0.95
closeAllPreboundMethod · 0.95
CloseMethod · 0.65