| 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 | |
| 323 | func (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"` |