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

Method doHealthCheck

duckdbservice/flight_handler.go:313–442  ·  view source on GitHub ↗
(body []byte, stream flight.FlightService_DoActionServer)

Source from the content-addressed store, hash-verified

311 if err := h.pool.validateControlMetadata(req.WorkerControlMetadata); err != nil {
312 return status.Errorf(codes.FailedPrecondition, "stale worker owner: %v", err)
313 }
314
315 if err := h.pool.DestroySession(req.SessionToken); err != nil {
316 return status.Errorf(codes.NotFound, "%v", err)
317 }
318
319 resp, _ := json.Marshal(map[string]bool{"ok": true})
320 return sendActionResult(stream, &flight.Result{Body: resp})
321}
322
323func (h *FlightSQLHandler) doHealthCheck(body []byte, stream flight.FlightService_DoActionServer) error {
324 var req server.WorkerHealthCheckPayload
325 if err := json.Unmarshal(body, &req); err != nil {
326 return status.Errorf(codes.InvalidArgument, "invalid HealthCheck request: %v", err)
327 }
328 if err := h.pool.validateControlMetadata(req.WorkerControlMetadata); err != nil {
329 return status.Errorf(codes.FailedPrecondition, "stale worker owner: %v", err)
330 }
331
332 // Block until warmup (extension loading + DuckLake attachment) completes.
333 // Without this, the control plane's waitForWorkerTCP health check passes
334 // as soon as the gRPC server starts, and clients get routed to a worker
335 // that hasn't attached DuckLake yet.
336 <-h.pool.warmupDone
337
338 // Kick a SELECT 1 liveness probe. Before this, the health check never
339 // executed SQL — it only read progress counters — so a DuckDB instance
340 // invalidated by an Internal/Fatal engine error passed every check and
341 // stayed schedulable, and the org's next connection was handed the dead
342 // instance. The probe runs asynchronously and we report the flag it sets;
343 // see probeInstanceLivenessAsync for why it must not block this response.
344 h.pool.probeInstanceLivenessAsync()
345 instanceInvalidated := h.pool.InstanceInvalidated()
346
347 // Poll DuckDB query progress for each active session.
348 //
349 // QueryProgress is a CGO call into DuckDB that *should* return instantly
350 // (it reads atomic progress counters), but can block if DuckDB holds an
351 // internal lock — for example, when the httpfs extension is mid-download
352 // on a large remote parquet file. If that happens while we hold the pool
353 // RLock, both the health check and any session create/close operations
354 // stall, the CP's 3-second health check timeout fires, and after 3
355 // consecutive failures the CP kills the worker — even though it's alive
356 // and making progress on the download.
357 //
358 // To prevent this:
359 // 1. Snapshot the session data we need under RLock, then release it.
360 // 2. Call QueryProgress outside the lock, with a per-session timeout.
361 // If the CGO call doesn't return within queryProgressTimeout, we
362 // report the session as "busy" (pct=-1) and skip stall detection
363 // for this cycle. The health check always responds promptly.
364 //
365 // Real crashes (process death) are detected by the K8s pod informer
366 // independently of health checks, so skipping stall detection during
367 // I/O-heavy operations doesn't create a blind spot for crash recovery.
368 type sessionProgressInfo struct {
369 Pct float64 `json:"pct"`
370 Rows uint64 `json:"rows"`

Calls 6

sendActionResultFunction · 0.85
AddMethod · 0.80
IsDrainingMethod · 0.80
ActiveSessionsMethod · 0.80
ActiveDrainWorkMethod · 0.80