ShutdownAll stops all workers gracefully.
()
| 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. |
| 1064 | func (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. |
| 1071 | func (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 { |