| 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. |
| 496 | func isWorkerConnPoolTimeoutError(err error) bool { |
| 497 | return err != nil && strings.Contains(err.Error(), "failed to obtain connection from pool") |
| 498 | } |
| 499 | |
| 500 | func (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 | |
| 513 | func (sm *SessionManager) beginSessionCreation(ctx context.Context) (context.Context, func() bool, error) { |
| 514 | return sm.lifecycle.begin(ctx) |
| 515 | } |
| 516 | |
| 517 | func (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 { |