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

Method acquireConnectionSlot

controlplane/session_mgr.go:195–271  ·  view source on GitHub ↗
(ctx context.Context, pid int32, username string, protocol string, profile *WorkerProfile)

Source from the content-addressed store, hash-verified

193 vcpus := 1
194 if requestedVCPUs != nil {
195 var err error
196 vcpus, err = requestedVCPUs(profile)
197 if err != nil {
198 return nil, err
199 }
200 }
201 if vcpus <= 0 {
202 return nil, fmt.Errorf("requested vcpus must be positive, got %d", vcpus)
203 }
204 limits := func(user string) configstore.OrgResourceLimits {
205 if resourceLimits == nil {
206 return configstore.OrgResourceLimits{}
207 }
208 return resourceLimits(user)
209 }
210 lease, err := limiter.Acquire(ctx, connectionAdmissionRequest{
211 PID: pid,
212 Username: username,
213 Protocol: protocol,
214 RequestedVCPUs: vcpus,
215 }, limits)
216 if err != nil {
217 return nil, err
218 }
219 if sm.lifecycle.isClosed() {
220 sm.releaseConnectionSlot(lease)
221 return nil, ErrSessionManagerDraining
222 }
223 return lease, nil
224 }
225
226 sm.mu.Lock()
227 if sm.lifecycle.isClosed() {
228 sm.mu.Unlock()
229 return nil, ErrSessionManagerDraining
230 }
231 if sm.maxConnections <= 0 || sm.activeSlots < sm.maxConnections {
232 sm.activeSlots++
233 sm.mu.Unlock()
234 return nil, nil
235 }
236 waiter := &connectionWaiter{ready: make(chan struct{})}
237 sm.waiters = append(sm.waiters, waiter)
238 sm.mu.Unlock()
239
240 select {
241 case <-waiter.ready:
242 if waiter.err != nil {
243 return nil, waiter.err
244 }
245 return nil, nil
246 case <-ctx.Done():
247 sm.mu.Lock()
248 defer sm.mu.Unlock()
249 if waiter.granted {
250 return nil, nil
251 }
252 sm.removeWaiterLocked(waiter)

Calls 5

releaseConnectionSlotMethod · 0.95
removeWaiterLockedMethod · 0.95
isClosedMethod · 0.80
AcquireMethod · 0.65
ErrMethod · 0.65