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

Method OnWorkerCrash

controlplane/session_mgr.go:848–913  ·  view source on GitHub ↗

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))

Source from the content-addressed store, hash-verified

846func (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()

Calls 8

detachSessionLockedMethod · 0.95
CloseMethod · 0.95
releaseSessionLeaseMethod · 0.95
MarkDeadMethod · 0.80
NowMethod · 0.80
beginCleanupMethod · 0.80
RequestRebalanceMethod · 0.80
CloseMethod · 0.65