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

Method createSessionOnWorker

controlplane/session_mgr.go:540–633  ·  view source on GitHub ↗
(ctx context.Context, username string, pid int32, memoryLimit string, threads int, worker *ManagedWorker, protocol string, retireOnFailure bool, lease connectionLease)

Source from the content-addressed store, hash-verified

538 sm.log.Error("Failed to load user persistent secrets; session starts without them.",
539 "pid", pid, "worker", worker.ID, "user", username, "error", secretErr)
540 }
541 }
542
543 sessionToken, secretWarnings, err := worker.CreateSession(ctx, username, memoryLimit, threads, secretStatements, pid)
544 if err != nil {
545 sm.log.Warn("Failed to create session on worker.",
546 "pid", pid,
547 "worker", worker.ID,
548 "user", username,
549 "protocol", protocol,
550 "duration", time.Since(createStart),
551 "retire_on_failure", retireOnFailure,
552 "error", err,
553 )
554 if retireOnFailure {
555 sm.pool.RetireWorkerIfNoSessions(worker.ID)
556 }
557 return 0, nil, fmt.Errorf("create session on worker %d: %w", worker.ID, err)
558 }
559
560 for _, w := range secretWarnings {
561 sm.log.Warn("User persistent secret replay warning.",
562 "pid", pid, "worker", worker.ID, "user", username, "warning", w)
563 }
564
565 executor := flightclient.NewFlightExecutorFromClientWithQueryLogLimiter(
566 worker.client,
567 sessionToken,
568 worker.workerQueryLogLimiter(),
569 )
570 executor.SetControlMetadata(worker.ID, worker.OwnerCPInstanceID(), worker.OwnerEpoch())
571
572 if pid == 0 {
573 pid = reservePID(globalNextPID)
574 }
575
576 session := &ManagedSession{
577 PID: pid,
578 Username: username,
579 WorkerID: worker.ID,
580 Protocol: protocol,
581 StartedAt: time.Now().UTC(),
582 SessionToken: sessionToken,
583 Executor: executor,
584 lease: lease,
585 }
586
587 sm.mu.Lock()
588 if sm.lifecycle.isClosed() {
589 sm.mu.Unlock()
590 sm.cleanupUnregisteredWorkerSession(worker, session)
591 return 0, nil, ErrSessionManagerDraining
592 }
593 sm.sessions[pid] = session
594 sm.byWorker[worker.ID] = append(sm.byWorker[worker.ID], pid)
595 sessionCount := len(sm.sessions)
596 workerSessionCount := len(sm.byWorker[worker.ID])
597 sm.mu.Unlock()

Callers 2

Calls 12

reservePIDFunction · 0.85
NowMethod · 0.80
SetControlMetadataMethod · 0.80
isClosedMethod · 0.80
RequestRebalanceMethod · 0.80
CreateSessionMethod · 0.65
OwnerCPInstanceIDMethod · 0.45
OwnerEpochMethod · 0.45
ErrorMethod · 0.45

Tested by

no test coverage detected