Custom action handlers (called via customActionServer.DoAction)
(body []byte, stream flight.FlightService_DoActionServer)
| 209 | // Custom action handlers (called via customActionServer.DoAction) |
| 210 | |
| 211 | func (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, |