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

Method doCreateSession

duckdbservice/flight_handler.go:211–263  ·  view source on GitHub ↗

Custom action handlers (called via customActionServer.DoAction)

(body []byte, stream flight.FlightService_DoActionServer)

Source from the content-addressed store, hash-verified

209// Custom action handlers (called via customActionServer.DoAction)
210
211func (h *FlightSQLHandler) doCreateSession(body []byte, stream flight.FlightService_DoActionServer) error {
212 var req server.WorkerCreateSessionPayload
213 if err := json.Unmarshal(body, &req); err != nil {
214 return status.Errorf(codes.InvalidArgument, "invalid CreateSession request: %v", err)
215 }
216 if req.Username == "" {
217 return status.Error(codes.InvalidArgument, "username is required")
218 }
219 if h.pool.IsDraining() {
220 return status.Error(codes.Unavailable, "worker is draining")
221 }
222 if err := h.pool.validateControlMetadata(req.WorkerControlMetadata); err != nil {
223 return status.Errorf(codes.FailedPrecondition, "stale worker owner: %v", err)
224 }
225
226 if h.pool.sharedWarmMode {
227 if _, err := h.pool.currentSessionConfig(); err != nil {
228 return status.Error(codes.FailedPrecondition, "worker is not activated")
229 }
230 }
231
232 // Validate username against configured users
233 if !h.pool.sharedWarmMode {
234 if _, ok := h.pool.cfg.Users[req.Username]; !ok {
235 return status.Error(codes.PermissionDenied, "unknown username")
236 }
237 }
238
239 session, secretWarnings, err := h.pool.CreateSession(req.Username, req.MemoryLimit, req.Threads, req.SecretStatements)
240 attachSessionLog(session, req.Username, req.PID)
241 if drainErr := workerDrainingStatus(err); drainErr != nil {
242 return drainErr
243 }
244 // Session create runs DDL against the instance, so it is where an ALREADY
245 // invalidated worker announces itself: the observed failure sequence is a
246 // fatal on one session, then the CP handing the same worker a brand new
247 // session ~2 minutes later that dies initializing its database metadata.
248 // Flagging here retires the worker on that first rejected create instead of
249 // waiting for the liveness probe. Opaque because this path replays the
250 // user's persistent CREATE SECRET statements: its error can echo a
251 // credential with no single statement to classify against.
252 h.pool.noteInstanceErrorOpaque(err)
253 if err != nil {
254 return status.Errorf(codes.ResourceExhausted, "create session: %v", err)
255 }
256
257 // secret_warnings is always present (empty slice, never null/absent): its
258 // presence is how the control plane distinguishes a worker that processed
259 // the secret replay from an old-image worker that silently ignored the
260 // payload field.
261 if secretWarnings == nil {
262 secretWarnings = []string{}
263 }
264 resp, _ := json.Marshal(map[string]any{
265 "session_token": session.ID,
266 "secret_warnings": secretWarnings,

Calls 8

workerDrainingStatusFunction · 0.85
sendActionResultFunction · 0.85
IsDrainingMethod · 0.80
currentSessionConfigMethod · 0.80
CreateSessionMethod · 0.65
DestroySessionMethod · 0.65
ErrorMethod · 0.45