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

Method CreateSessionWithProtocol

controlplane/session_mgr.go:334–449  ·  view source on GitHub ↗
(ctx context.Context, username string, pid int32, memoryLimit string, threads int, protocol string, profile *WorkerProfile)

Source from the content-addressed store, hash-verified

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),

Calls 15

beginSessionCreationMethod · 0.95
ReservePIDMethod · 0.95
acquireConnectionSlotMethod · 0.95
releaseConnectionSlotMethod · 0.95
resolveSessionLimitsMethod · 0.95
missReasonMethod · 0.95
createSessionOnWorkerMethod · 0.95
capacityRetrySecondsFunction · 0.85
isWorkerSessionCapErrorFunction · 0.85