OnWorkerCrash handles a worker crash by marking all affected executors as dead and notifying sessions. Executors are marked dead BEFORE the shared gRPC client is closed to prevent nil-pointer panics from concurrent RPCs. errorFn is called for each affected session to send an error to the client.
(workerID int, errorFn func(pid int32))
| 846 | func (sm *SessionManager) OnWorkerCrash(workerID int, errorFn func(pid int32)) { |
| 847 | sm.mu.Lock() |
| 848 | pids := make([]int32, len(sm.byWorker[workerID])) |
| 849 | copy(pids, sm.byWorker[workerID]) |
| 850 | |
| 851 | // Mark all executors as dead first (under lock) so any concurrent RPC |
| 852 | // sees the dead flag before the gRPC client is closed. |
| 853 | for _, pid := range pids { |
| 854 | if s, ok := sm.sessions[pid]; ok && s.Executor != nil { |
| 855 | s.Executor.MarkDead() |
| 856 | } |
| 857 | } |
| 858 | sm.mu.Unlock() |
| 859 | |
| 860 | sm.log.Warn("Worker crashed, notifying sessions.", "worker", workerID, "sessions", len(pids), "pids", pids) |
| 861 | |
| 862 | for _, pid := range pids { |
| 863 | cleanupStart := time.Now() |
| 864 | errorFn(pid) |
| 865 | var executor *flightclient.FlightExecutor |
| 866 | var connCloser io.Closer |
| 867 | var session *ManagedSession |
| 868 | var finishCleanup func() |
| 869 | sm.mu.Lock() |
| 870 | detached, remainingSessions, _, ok := sm.detachSessionLocked(pid) |
| 871 | if ok { |
| 872 | session = detached |
| 873 | executor = detached.Executor |
| 874 | connCloser = detached.connCloser |
| 875 | finishCleanup = sm.lifecycle.beginCleanup() |
| 876 | } |
| 877 | sm.mu.Unlock() |
| 878 | if executor != nil { |
| 879 | _ = executor.Close() |
| 880 | } |
| 881 | // Close the TCP connection to unblock the message loop's read. |
| 882 | // This causes the session goroutine to exit instead of looping |
| 883 | // with ErrWorkerDead on every query. The deferred close in |
| 884 | // handleConnection will also call Close() on the same conn; |
| 885 | // that's harmless (net.Conn.Close on a closed socket returns |
| 886 | // an error which is discarded). |
| 887 | if connCloser != nil { |
| 888 | _ = connCloser.Close() |
| 889 | } |
| 890 | sm.releaseSessionLease(session, "pid", pid) |
| 891 | sm.log.Info("Worker crash session cleanup completed.", |
| 892 | "pid", pid, |
| 893 | "worker", workerID, |
| 894 | "session_found", ok, |
| 895 | "duration", time.Since(cleanupStart), |
| 896 | "remaining_sessions", remainingSessions, |
| 897 | ) |
| 898 | if ok { |
| 899 | finishCleanup() |
| 900 | } |
| 901 | } |
| 902 | |
| 903 | sm.mu.Lock() |
| 904 | delete(sm.byWorker, workerID) |
| 905 | sm.mu.Unlock() |