| 332 | // publishes the session. Reclaim that raced success before returning |
| 333 | // it to a client that shutdown has already rejected. |
| 334 | sm.DestroySession(resultPID) |
| 335 | resultPID = 0 |
| 336 | resultExecutor = nil |
| 337 | } |
| 338 | resultErr = ErrSessionManagerDraining |
| 339 | }() |
| 340 | |
| 341 | lease, err := sm.acquireConnectionSlot(ctx, pid, username, protocol, profile) |
| 342 | if err != nil { |
| 343 | return 0, nil, err |
| 344 | } |
| 345 | success := false |
| 346 | defer func() { |
| 347 | if !success { |
| 348 | sm.releaseConnectionSlot(lease) |
| 349 | } |
| 350 | }() |
| 351 | |
| 352 | memoryLimit, threads = sm.resolveSessionLimits(memoryLimit, threads) |
| 353 | |
| 354 | // Acquire a worker. Backend implementations may reuse warm workers, queue, |
| 355 | // spawn, or return a typed capacity error when no worker is immediately available. |
| 356 | observeControlPlaneWorkerQueueDepthDelta(1) |
| 357 | defer observeControlPlaneWorkerQueueDepthDelta(-1) |
| 358 | |
| 359 | // Acquire a worker and create the session on it. Normally one pass. If the |
| 360 | // worker rejects our session because it already holds its max session — a |
| 361 | // CP↔worker accounting drift that must never happen under one-session-per- |
| 362 | // worker — we do NOT fail the client for our own broken logic: recycle the |
| 363 | // inconsistent worker and try a fresh one (bounded), logging loudly so the |
| 364 | // drift is visible. ctx is the budget; each attempt also re-checks it. |
| 365 | var lastCapDriftErr error |
| 366 | for attempt := 1; attempt <= maxWorkerSessionCapDriftRetries+1; attempt++ { |
| 367 | if err := ctx.Err(); err != nil { |
| 368 | return 0, nil, err |
| 369 | } |
| 370 | |
| 371 | acquireStart := time.Now() |
| 372 | actx, acquireSpan := server.Tracer().Start(ctx, "duckgres.worker_acquire") |
| 373 | sm.log.Debug("Acquiring worker for session.", "pid", pid, "user", username, "attempt", attempt) |
| 374 | worker, err := sm.pool.AcquireWorker(actx, profile) |
| 375 | if err != nil { |
| 376 | var capacityErr *WorkerCapacityExhaustedError |
| 377 | if errors.As(err, &capacityErr) { |
| 378 | missReason := capacityErr.missReason() |
| 379 | observeControlPlaneWorkerAcquireFailure("worker_capacity_exhausted") |
| 380 | observeControlPlaneWorkerAcquireFailure("worker_capacity_" + string(missReason)) |
| 381 | acquireSpan.SetAttributes( |
| 382 | attribute.String("worker_capacity.reason", string(missReason)), |
| 383 | attribute.Int("worker_capacity.retry_after_seconds", capacityRetrySeconds(capacityErr.RetryAfter)), |
| 384 | ) |
| 385 | sm.log.Warn("Worker acquisition failed.", |
| 386 | "pid", pid, |
| 387 | "user", username, |
| 388 | "duration", time.Since(acquireStart), |
| 389 | "reason", missReason, |
| 390 | "retry_after", capacityErr.RetryAfter, |
| 391 | "retry_after_seconds", capacityRetrySeconds(capacityErr.RetryAfter), |