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

Method ReconnectFlightSession

controlplane/session_mgr.go:489–534  ·  view source on GitHub ↗
(ctx context.Context, username string, workerID int, ownerEpoch int64)

Source from the content-addressed store, hash-verified

487}
488
489// isWorkerConnPoolTimeoutError reports whether err is the worker-side
490// "failed to obtain connection from pool" session-create failure
491// (duckdbservice acquires the single-session DB connection with a 30s
492// timeout; see the MaxOpenConns=1 isolation contract). A worker returning
493// this is wedged — its connection never came back from the previous
494// session's cleanup — and because it parks hot-idle afterwards, plain
495// client retries deterministically reuse it; it must be recycled instead.
496func isWorkerConnPoolTimeoutError(err error) bool {
497 return err != nil && strings.Contains(err.Error(), "failed to obtain connection from pool")
498}
499
500func (sm *SessionManager) resolveSessionLimits(memoryLimit string, threads int) (string, int) {
501 if sm.rebalancer == nil {
502 return memoryLimit, threads
503 }
504 if memoryLimit == "" {
505 memoryLimit = sm.rebalancer.MemoryLimit()
506 }
507 if threads <= 0 {
508 threads = sm.rebalancer.PerSessionThreads()
509 }
510 return memoryLimit, threads
511}
512
513func (sm *SessionManager) beginSessionCreation(ctx context.Context) (context.Context, func() bool, error) {
514 return sm.lifecycle.begin(ctx)
515}
516
517func (sm *SessionManager) createSessionOnWorker(ctx context.Context, username string, pid int32, memoryLimit string, threads int, worker *ManagedWorker, protocol string, retireOnFailure bool, lease connectionLease) (int32, *flightclient.FlightExecutor, error) {
518 createStart := time.Now()
519 sm.log.Info("Creating session on worker.",
520 "pid", pid,
521 "worker", worker.ID,
522 "user", username,
523 "protocol", protocol,
524 "memory_limit", memoryLimit,
525 "threads", threads,
526 "owner_cp_instance_id", worker.OwnerCPInstanceID(),
527 "owner_epoch", worker.OwnerEpoch(),
528 )
529 // Load the user's persistent secrets for replay. Failure to load degrades
530 // to a session without user secrets (logged loudly) rather than a refused
531 // connection: the config store being briefly unavailable must not lock
532 // every returning user out of their warehouse.
533 var secretStatements []string
534 if sm.userSecretLoader != nil {
535 var secretErr error
536 secretStatements, secretErr = sm.userSecretLoader(ctx, username)
537 if secretErr != nil {

Calls 8

beginSessionCreationMethod · 0.95
ReservePIDMethod · 0.95
acquireConnectionSlotMethod · 0.95
releaseConnectionSlotMethod · 0.95
createSessionOnWorkerMethod · 0.95
ReconnectFlightWorkerMethod · 0.65
ReleaseWorkerMethod · 0.65